Skip to main content

sc_network_statement/
lib.rs

1// This file is part of Substrate.
2
3// Copyright (C) Parity Technologies (UK) Ltd.
4// SPDX-License-Identifier: GPL-3.0-or-later WITH Classpath-exception-2.0
5
6// This program is free software: you can redistribute it and/or modify
7// it under the terms of the GNU General Public License as published by
8// the Free Software Foundation, either version 3 of the License, or
9// (at your option) any later version.
10
11// This program is distributed in the hope that it will be useful,
12// but WITHOUT ANY WARRANTY; without even the implied warranty of
13// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
14// GNU General Public License for more details.
15
16// You should have received a copy of the GNU General Public License
17// along with this program. If not, see <https://www.gnu.org/licenses/>.
18
19//! Statement handling to plug on top of the network service.
20//!
21//! This crate implements gossip-based propagation of statements between nodes, layered on the
22//! substrate notifications protocol. Two protocol versions are negotiated per peer: `statement/2`
23//! (preferred) with `statement/1` as a fallback.
24//!
25//! ## Propagation
26//!
27//! - During major chain synchronization, statement gossip is paused so peers prioritize downloading
28//!   blocks; it resumes automatically once the node is fully synced (peers are reconnected to
29//!   recover statements missed while syncing).
30//! - A propagation loop runs every second (`config::PROPAGATE_TIMEOUT`): it takes all statements
31//!   added since the previous round and queues their hashes to a per-peer outbox. Each peer has at
32//!   most one propagation chunk in flight at a time. When its send slot is free, statements are
33//!   fetched from the store, filtered and encoded up to the maximum notification size
34//!   (`config::MAX_STATEMENT_NOTIFICATION_SIZE`, ~1 MiB).
35//! - Incoming statements are pushed onto a bounded validation queue
36//!   (`config::MAX_PENDING_STATEMENTS`); if the queue is full, incoming statements are dropped.
37//! - Peer reputation is adjusted based on statement quality (good, duplicate, invalid, flooding).
38//!
39//! ## Tracking received statements
40//!
41//! There is no per-peer record of delivered statements. A statement is propagated once, on the
42//! tick after its import: the propagation pass drains the store's recent set, so no later pass
43//! can pick it up again. The only duplicates worth preventing are statements sent back to the
44//! peers they came from.
45//!
46//! While a statement waits for validation, the peers it came from are recorded in
47//! `pending_statements_peers`. On import they move to `recently_received_statements`, and a
48//! peer that resends the statement before its tick is added there too. Propagation and
49//! initial-sync chunks both skip the recorded peers, and the propagation pass clears
50//! `recently_received_statements` when done.
51//!
52//! ## Initial sync
53//!
54//! A peer joining (or changing its topic affinity) receives the store's existing statements
55//! through a cursor over the store's admission journal: scheduling captures the journal's
56//! watermark, and bursts walk the admissions below it chunk by chunk, advancing the cursor as
57//! each send is confirmed — a failed chunk is resent from the same position. The watermark also
58//! splits delivery ownership between the two paths: admissions below the peer's highest
59//! watermark are the sync cursor's job, and propagation skips them, so within one sync a
60//! statement reaches the peer through at most one path. Duplicates stay possible at the edges —
61//! an affinity-change re-sync restarts the cursor from zero and redelivers earlier admissions,
62//! and a chunk whose send timed out may have arrived regardless, so its resend repeats it.
63//!
64//! ## Send scheduling
65//!
66//! Every peer has one send slot, shared by propagation and initial sync, so at most one chunk per
67//! connection is in flight. Each chunk carries a fresh id, so the result of a send left over from
68//! a previous connection cannot free the current one's slot. A completed chunk frees the slot at
69//! once and the next propagation chunk follows, a failed send included: a failed propagation
70//! chunk is not retried, but the rest of the backlog keeps draining, while a failed initial-sync
71//! chunk is resent from the sync's cursor.
72//!
73//! Propagation queues hashes in a per-peer outbox and fetches, filters and encodes them only when
74//! the slot is free, so a slow peer holds one encoded chunk rather than its whole backlog. An
75//! outbox past `config::MAX_PROPAGATION_OUTBOX_LEN` drops its oldest hashes, since the freshest
76//! statements are the ones still worth delivering. Dropped hashes are counted in
77//! `undelivered_statements`.
78//!
79//! In-flight bytes of both kinds are held against the shared
80//! `config::MAX_SEND_IN_FLIGHT_BYTES` budget. A peer that finds the budget full, or whose store
81//! fetch fails, is parked once and refilled in parking order as completed sends free bytes and on
82//! propagation ticks. While initial syncs are pending, propagation parks
83//! `config::INITIAL_SYNC_RESERVED_BYTES` early: refills reclaim freed bytes synchronously,
84//! while the timer-driven sync bursts would otherwise always find the budget full.
85//!
86//! ## Topic affinity and light nodes
87//!
88//! The `statement/2` protocol lets a peer advertise which topics it cares about as a bloom filter
89//! ("topic affinity"). Once a peer has an active affinity filter, only matching statements are
90//! forwarded to it; when its affinity changes, newly relevant statements are re-sent. Affinity
91//! advertisements are rate-limited. See the `affinity` module.
92//!
93//! Light-client peers on `statement/2` must advertise an affinity before receiving any statements:
94//! a light V2 peer pulls only the topics it cares about instead of the full feed, and is synced
95//! those statements in an initial burst on connect. Full nodes receive all statements unless they
96//! opt into an affinity.
97//!
98//! ## Usage
99//!
100//! - Use [`StatementHandlerPrototype::new`] to create a prototype.
101//! - Pass the `NonDefaultSetConfig` returned from [`StatementHandlerPrototype::new`] to the network
102//!   configuration as an extra peers set.
103//! - Use [`StatementHandlerPrototype::build`] then [`StatementHandler::run`] to obtain a `Future`
104//!   that processes statements.
105
106mod 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
161/// A set of statements.
162pub type Statements = Vec<Statement>;
163
164type StatementBatch = BoundedVec<Statement, ConstU32<{ MAX_STATEMENTS_PER_NOTIFICATION as u32 }>>;
165
166/// The protocol version that was negotiated with a peer.
167#[derive(Debug, Clone, Copy, PartialEq, Eq)]
168enum PeerProtocolVersion {
169	/// V1: messages are encoded as `Vec<Statement>` (the legacy format).
170	V1,
171	/// V2: messages are encoded as `StatementMessage` enum (supports topic affinity).
172	V2,
173}
174
175impl PeerProtocolVersion {
176	/// Returns the encoding envelope overhead for this protocol version.
177	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	/// Bloom filter bytes representing the topics this peer is interested in.
190	#[codec(index = 1)]
191	ExplicitTopicAffinity(AffinityFilter),
192}
193
194/// Codec variant index for `StatementMessage::Statements`, kept in sync with `#[codec(index)]`.
195const STATEMENTS_VARIANT_INDEX: u8 = 0;
196
197impl StatementMessage {
198	/// Encode a slice of statement references as a `StatementMessage::Statements`
199	/// without cloning the statements.
200	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
208/// Future resolving to statement import result.
209pub type StatementImportFuture = oneshot::Receiver<SubmitResult>;
210
211mod rep {
212	use sc_network::ReputationChange as Rep;
213	/// Reputation change when a peer sends us any statement.
214	///
215	/// This forces node to verify it, thus the negative value here. Once statement is verified,
216	/// reputation change should be refunded with `ANY_STATEMENT_REFUND`
217	pub const ANY_STATEMENT: Rep = Rep::new(-(1 << 4), "Any statement");
218	/// Reputation change when a peer sends us any statement that is not invalid.
219	pub const ANY_STATEMENT_REFUND: Rep = Rep::new(1 << 4, "Any statement (refund)");
220	/// Reputation change when a peer sends us an statement that we didn't know about.
221	pub const GOOD_STATEMENT: Rep = Rep::new(1 << 8, "Good statement");
222	/// Reputation change when a peer sends us an invalid statement.
223	pub const INVALID_STATEMENT: Rep = Rep::new(-(1 << 12), "Invalid statement");
224	/// Reputation change when a peer sends us a duplicate statement.
225	pub const DUPLICATE_STATEMENT: Rep = Rep::new(-(1 << 7), "Duplicate statement");
226	/// Reputation change when a peer floods us with statements.
227	pub const STATEMENT_FLOODING: Rep = Rep::new_fatal("Statement flooding");
228	/// Reputation change when a peer sends us a message we can't decode.
229	pub const BAD_MESSAGE: Rep = Rep::new(-(1 << 12), "Bad statement message");
230}
231
232const LOG_TARGET: &str = "statement-gossip";
233/// V2 statement protocol suffix, work in progress protocol with topic affinity and other
234/// improvements, may have breaking changes before stabilization.
235const STATEMENT_PROTOCOL_V2: &str = "statement/2";
236/// V1 statement protocol suffix, current stable protocol, no breaking changes will be made to it.
237const STATEMENT_PROTOCOL_V1: &str = "statement/1";
238/// Maximum time we wait for sending a notification to a peer.
239const SEND_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(10);
240/// Interval for sending statement batches during initial sync to new peers.
241const INITIAL_SYNC_BURST_INTERVAL: std::time::Duration = std::time::Duration::from_millis(10);
242
243/// Maximum admission-journal entries one initial-sync fetch visits, so that a peer whose
244/// affinity matches nothing cannot make one burst walk the whole journal in the event loop.
245const INITIAL_SYNC_SCAN_LIMIT: usize = 4096;
246/// Interval for processing pending topic affinity changes from peers.
247const PENDING_AFFINITIES_INTERVAL: std::time::Duration = std::time::Duration::from_secs(1);
248/// Delay before re-adding a peer to the reserved set after a forced disconnect for sync recovery.
249const 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	/// The network layer rejected the send.
276	pub const NETWORK: &str = "network";
277	/// The send did not complete within `SEND_TIMEOUT`.
278	pub const TIMEOUT: &str = "timeout";
279	/// The chunk was never handed to the network because the peer had no message sink.
280	pub const NO_SINK: &str = "no_sink";
281	/// The peer's propagation outbox overflowed and the oldest queued hashes were dropped.
282	pub const OUTBOX_FULL: &str = "outbox_full";
283}
284
285mod sync_outcome {
286	/// A burst found the backlog drained, so every statement reached the peer.
287	pub const COMPLETED: &str = "completed";
288	/// The sync ended before its backlog was drained.
289	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 from 1μs to ~1s covering microsecond to millisecond range.
384					.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
465/// Prototype for a [`StatementHandler`].
466pub struct StatementHandlerPrototype {
467	protocol_name: ProtocolName,
468	notification_service: Box<dyn NotificationService>,
469}
470
471impl StatementHandlerPrototype {
472	/// Create a new instance.
473	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	/// Turns the prototype into the actual handler.
512	///
513	/// Important: the statements handler is initially disabled and doesn't gossip statements.
514	/// Gossiping is enabled when major syncing is done.
515	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		// Still bounded via the `MAX_PENDING_STATEMENTS` check in `on_statements`.
530		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
623/// Handler for statements. Call [`StatementHandler::run`] to start the processing.
624pub struct StatementHandler<
625	N: NetworkPeers + NetworkEventStream,
626	S: SyncEventStream + sp_consensus::SyncOracle,
627> {
628	protocol_name: ProtocolName,
629	/// Interval at which we call `propagate_statements`.
630	propagate_timeout: stream::Fuse<Pin<Box<dyn Stream<Item = ()> + Send>>>,
631	/// Pending statements verification tasks.
632	pending_statements:
633		FuturesUnordered<Pin<Box<dyn Future<Output = (Hash, Option<SubmitResult>)> + Send>>>,
634	/// As multiple peers can send us the same statement, we group
635	/// these peers using the statement hash while the statement is
636	/// imported. This prevents that we import the same statement
637	/// multiple times concurrently.
638	pending_statements_peers: HashMap<Hash, HashSet<PeerId>>,
639	/// Statements received from peers and imported since the last propagation
640	/// pass, each with the peers that sent it. Propagation skips those peers,
641	/// so a statement never returns to a peer it came from. Cleared after each
642	/// pass.
643	recently_received_statements: HashMap<Hash, HashSet<PeerId>>,
644	/// Network service to use to send messages and manage peers.
645	network: N,
646	/// Syncing service.
647	sync: S,
648	/// Receiver for syncing-related events.
649	sync_event_stream: stream::Fuse<Pin<Box<dyn Stream<Item = SyncEvent> + Send>>>,
650	/// Notification service.
651	notification_service: Box<dyn NotificationService>,
652	// All connected peers
653	peers: HashMap<PeerId, Peer>,
654	statement_store: Arc<dyn StatementStore>,
655	queue_sender: async_channel::Sender<(Statement, oneshot::Sender<SubmitResult>)>,
656	/// Maximum statements per second per peer.
657	statements_per_second: NonZeroU32,
658	/// Prometheus metrics.
659	metrics: Option<Metrics>,
660	/// Timeout for sending next statement batch during initial sync.
661	initial_sync_timeout: Pin<Box<dyn FusedFuture<Output = ()> + Send>>,
662	/// Timeout for processing pending topic affinity changes.
663	pending_affinities_timeout: Pin<Box<dyn FusedFuture<Output = ()> + Send>>,
664	/// Pending initial syncs per peer.
665	pending_initial_syncs: HashMap<PeerId, PendingInitialSync>,
666	/// Queue for round-robin processing of initial syncs.
667	initial_sync_peer_queue: VecDeque<PeerId>,
668	/// Next value to hand out as [`PendingInitialSync::sync_id`].
669	next_initial_sync_id: u64,
670	/// Encoded bytes of initial-sync chunks in `pending_sends`, held against the shared
671	/// [`MAX_SEND_IN_FLIGHT_BYTES`] budget.
672	initial_sync_in_flight_bytes: u64,
673	/// Statement hashes queued for propagation to each peer, drained from the front as chunks
674	/// are fetched, whether the fetch yields a send or not.
675	propagation_outboxes: HashMap<PeerId, VecDeque<Hash>>,
676	/// Id of the chunk in flight, per peer — the peer's single send slot, shared by
677	/// propagation and initial sync.
678	in_flight_chunks: HashMap<PeerId, u64>,
679	/// Next value to hand out as [`PendingSendResult::chunk_id`].
680	next_chunk_id: u64,
681	/// Encoded bytes of propagation chunks in `pending_sends`, held against the shared
682	/// [`MAX_SEND_IN_FLIGHT_BYTES`] budget.
683	propagation_in_flight_bytes: u64,
684	/// Peers whose propagation chunk was deferred because the shared byte budget was
685	/// exhausted or the store fetch failed, refilled in parking order as bytes free up
686	/// and on propagation ticks. A peer parks at most once.
687	parked_propagations: VecDeque<PeerId>,
688	/// Pending propagation sends, polled by the main event loop.
689	pending_sends: PendingSends,
690	/// Tracks peers that connected while major sync was active and adds them to the reserved set
691	/// once sync ends
692	deferred_peers: HashSet<PeerId>,
693	/// Set to `true` when an incoming statement is dropped because `is_major_syncing()` is true
694	dropped_statements_during_sync: bool,
695	/// Peer scheduled for forced disconnect+reconnect to recover statements missed during sync
696	sync_recovery_peer: Option<PeerId>,
697	/// Fires when the `sync_recovery_peer` re-add delay has elapsed
698	sync_recovery_readd_timeout: Pin<Box<dyn FusedFuture<Output = ()> + Send>>,
699}
700
701/// A token bucket, measured against whatever clock it was built with.
702trait TokenBucket: fmt::Debug + Send + Sync {
703	/// Whether admitting `n` more cells would exceed the quota.
704	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/// Per-peer rate limiter using a token bucket algorithm.
718///
719/// The token bucket allows short bursts up to the per-second limit while enforcing
720/// the average rate over time.
721#[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	/// The same quota, measured against `clock`.
732	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	/// Check if receiving `count` statements would exceed the rate limit.
742	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/// Peer information
755#[cfg_attr(not(any(test, feature = "test-helpers")), doc(hidden))]
756#[derive(Debug)]
757pub struct Peer {
758	/// Rate limiter for statement flooding protection.
759	rate_limiter: PeerRateLimiter,
760	/// Protocol version negotiated with this peer.
761	protocol_version: PeerProtocolVersion,
762	/// Topic affinity filter received from a v2 peer.
763	/// When set, only statements matching this filter should be propagated to the peer.
764	topic_affinity: Option<AffinityFilter>,
765	/// Whether this peer is a light client.
766	/// Light clients on V2 must set topic affinity before receiving statements.
767	is_light: bool,
768	/// A pending topic affinity filter waiting to be scheduled for initial sync.
769	/// Set when a new `ExplicitTopicAffinity` arrives; consumed by the main loop
770	/// once any in-progress initial sync for this peer completes.
771	pending_topic_affinity: Option<AffinityFilter>,
772	/// One past the newest admission covered by the peer's initial syncs.
773	sync_watermark: u64,
774}
775
776/// Tracks pending initial sync state for a peer as a cursor over the store's admission
777/// journal (statements fetched on-demand).
778struct PendingInitialSync {
779	/// Admission sequence number the next burst resumes fetching from.
780	cursor: u64,
781	/// One past the newest admission covered by this sync.
782	watermark: u64,
783	started_at: Instant,
784	/// Identifies this scheduling, so that a chunk still in flight from a previous one can be told
785	/// apart once its result arrives.
786	sync_id: u64,
787}
788
789enum SendOutcome {
790	/// The notification was accepted by the network layer.
791	Sent,
792	/// The network layer rejected the send.
793	NetworkError(error::Error),
794	/// The send did not complete within `SEND_TIMEOUT`.
795	TimedOut,
796}
797
798enum SendKind {
799	Propagation,
800	InitialSync {
801		sync_id: u64,
802		/// Cursor position the pending sync advances to once this chunk's send is confirmed.
803		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
816/// Result of an asynchronous send.
817struct PendingSendResult {
818	peer: PeerId,
819	statement_count: usize,
820	bytes_sent: u64,
821	result: SendOutcome,
822	kind: SendKind,
823	/// Id of the sent chunk, propagation or initial-sync. The result is stale if the id
824	/// doesn't match the peer's send slot in `in_flight_chunks`.
825	chunk_id: u64,
826}
827
828/// Type alias for the pending sends future collection, this is a list of in-flight sends to peers.
829type PendingSends =
830	FuturesUnordered<Pin<Box<dyn Future<Output = PendingSendResult> + Send + 'static>>>;
831
832/// Encoding overhead for V1: just the `Compact<u32>` vec length prefix (max 5 bytes).
833const V1_ENVELOPE_OVERHEAD: usize = 5;
834
835/// Encoding overhead for V2: 1 byte enum discriminant + `Compact<u32>` vec length prefix.
836const V2_ENVELOPE_OVERHEAD: usize = 1 + V1_ENVELOPE_OVERHEAD;
837
838/// Returns the maximum payload size for statement notifications given the
839/// protocol envelope overhead.
840fn 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
856/// Fetch the next chunk of statements admitted between `cursor` and `watermark`, filtering
857/// in the `admitted_statements` callback so non-matching statements are never cloned into
858/// the batch.
859///
860/// Returns the batch and the accumulated encoded size. A size above `max_size` signals a
861/// lone oversized statement: it is taken and the cursor sits past it, so the caller must
862/// drop the chunk and move on.
863fn 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			// The peer supplied this statement, do not send it back.
887			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
901/// Fetch the next chunk of statements for a peer from `hashes`, filtering in the
902/// `statements_by_hashes` callback so non-matching statements are never materialized.
903fn 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			// The peer supplied this statement, do not send it back.
923			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
947/// Whether the peer sent us the statement, directly or while it was queued for
948/// validation. `pending_statements_peers` covers the race where a statement is
949/// drained for propagation while the peers that sent it still sit there.
950fn 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	/// Create a new peer for testing/benchmarking purposes.
962	#[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	/// Whether this peer is ready to receive statements.
975	///
976	/// Light V2 peers must set their topic affinity before receiving any statements.
977	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	/// Create a new `StatementHandler` for testing/benchmarking purposes.
998	#[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	/// Get mutable access to pending statements for testing/benchmarking.
1046	#[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	/// Turns the [`StatementHandler`] into a future that should run forever and not be
1055	/// interrupted.
1056	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						// Syncing has seemingly closed. Closing as well.
1077						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						// `Notifications` has seemingly closed. Closing as well.
1085						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	/// Release excess map capacity once the map is less than a quarter full.
1112	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	/// Record a send attempt that never reached the peer.
1121	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	/// Record a failed send whose statements' delivery to the peer is abandoned.
1128	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	/// Add all peers that were deferred during major sync to the reserved set
1139	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	/// Pick one connected peer, remove it from the reserved set (forcing a disconnect), and
1165	/// schedule it for re-adding after `SYNC_RECOVERY_READD_DELAY`. When the peer reconnects it
1166	/// performs a fresh initial sync, delivering any statements that were dropped while the
1167	/// `is_major_syncing` guard was active
1168	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	/// Re-adds the sync-recovery peer to the reserved set after the backoff window has elapsed
1201	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	/// React to peer connect/disconnect events from the sync subsystem:
1218	///
1219	/// - On connect while major-syncing: defer the peer (kept in `deferred_peers`) instead of
1220	///   adding it to the statement protocol's reserved set, prioritizing block download; deferred
1221	///   peers are flushed once syncing finishes.
1222	/// - On connect otherwise: add the peer to the reserved set.
1223	/// - On disconnect: remove the peer from the reserved set, or from the deferred set if it never
1224	///   joined.
1225	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	/// Dispatch a notification-protocol event for the statement protocol:
1262	///
1263	/// - Validates inbound substreams by peer role.
1264	/// - Tracks stream open/close to maintain per-peer state.
1265	/// - Decodes incoming notifications: V1 peers send raw statement batches; V2 peers send a
1266	///   `StatementMessage` that is either a batch of statements or an `ExplicitTopicAffinity`
1267	///   advertisement.
1268	/// - Rate-limits affinity advertisements (reporting `rep::BAD_MESSAGE` on abuse); otherwise
1269	///   stores the filter as pending until applied by the main loop.
1270	async fn handle_notification_event(&mut self, event: NotificationEvent) {
1271		match event {
1272			NotificationEvent::ValidateInboundSubstream { peer, handshake, result_tx, .. } => {
1273				// Only accept peers whose role can be determined
1274				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				// If negotiated_fallback is Some, the peer connected on a fallback protocol
1287				// (v1). If None, the peer connected on the main protocol (v2).
1288				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				// Light V2 peers must set topic affinity before receiving statements.
1332				// All other peers get initial sync immediately.
1333				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				// Accept statements only when node is not major syncing
1364				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						// V1 peers send raw Vec<Statement>.
1381						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						// V2 peers send StatementMessage enum.
1395						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											// Defer both the affinity update and sync scheduling
1414											// to the main loop tick.
1415											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	/// Handle a batch of statements received from a peer.
1434	///
1435	/// For the batch:
1436	/// - Enforces the per-peer rate limit — on abuse, disconnects the peer and reports
1437	///   `rep::STATEMENT_FLOODING`.
1438	/// - Skips statements already in the store, reporting `rep::DUPLICATE_STATEMENT` if the same
1439	///   peer sent it twice.
1440	/// - Enqueues unknown statements onto the bounded validation queue.
1441	/// - Drops the remaining statements in the batch if the queue is full
1442	///   (`MAX_PENDING_STATEMENTS`).
1443	#[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				// Initiate peer state cleanup in the `NotificationStreamClosed` handler
1463				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 the statement still awaits its propagation pass, record the
1495					// peer so the pass does not send it back. Only join an existing
1496					// entry, or replays would grow the map without bound. Peers that
1497					// sent a not yet imported statement are tracked in
1498					// `pending_statements_peers` and move here on import.
1499					if let Some(peers) = self.recently_received_statements.get_mut(&hash) {
1500						peers.insert(who);
1501					}
1502
1503					// The store can already hold the statement while its validation
1504					// completion still awaits processing by the event loop. Join the
1505					// pending entry so the peer moves to `recently_received_statements`
1506					// on import, exactly as if the message had arrived before the
1507					// store insert.
1508					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							// Already received this from the same peer.
1555							self.network.report_peer(who, rep::DUPLICATE_STATEMENT);
1556						}
1557					},
1558				}
1559
1560				statements_left -= 1;
1561			}
1562		}
1563	}
1564
1565	/// Adjust the sending peer's reputation based on the outcome of importing a statement it sent.
1566	///
1567	/// Every newly received statement is first charged `rep::ANY_STATEMENT` (a small **decrease**)
1568	/// in [`on_statements`](Self::on_statements); this method applies the follow-up adjustment
1569	/// once the statement has been validated:
1570	///
1571	/// - `New` → **increase** by `rep::GOOD_STATEMENT` — a valid, previously unknown statement; the
1572	///   net change is positive (the reward outweighs the initial charge).
1573	/// - `Known` → **increase** by `rep::ANY_STATEMENT_REFUND`, which exactly cancels the initial
1574	///   `rep::ANY_STATEMENT` charge (net zero) — valid but already in the store.
1575	/// - `Invalid` → **decrease** by `rep::INVALID_STATEMENT`, a large penalty — the statement
1576	///   failed validation.
1577	/// - `KnownExpired`, `Rejected`, `InternalError` → no follow-up change, so the peer keeps the
1578	///   initial `rep::ANY_STATEMENT` charge.
1579	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	/// Handle a completed validation task. Adjusts the reputation of every peer
1591	/// that sent us the statement and, if the statement awaits propagation,
1592	/// records those peers so it is not sent back to them.
1593	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				// `New` and `Known` mean the statement awaits the next propagation
1600				// pass. Remember who sent it so the pass does not send it back.
1601				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	/// Queue the given `statements` for propagation to the given `peer`.
1611	///
1612	/// Internally filters out statements the peer sent to us.
1613	/// For v2 peers with a topic affinity filter, also filters by topic match.
1614	/// Surviving hashes are appended to the peer's outbox.
1615	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			// The peer supplied this statement, do not send it back.
1636			if has_received_from(recently_received_statements, pending_statements_peers, hash, who)
1637			{
1638				return None;
1639			}
1640			// For v2 peers with topic affinity, filter by topic match.
1641			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			// The freshest statements are the ones still worth delivering, so an
1652			// overflowing outbox drops from the front.
1653			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	/// Send the next propagation chunk to `who` if its send slot is free.
1670	///
1671	/// Statements are fetched, filtered and encoded only when a chunk actually goes out, so a
1672	/// slow peer holds one encoded chunk, not its whole backlog. Hashes whose statements left
1673	/// the store since they were queued are dropped.
1674	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			// Admission is checked before fetching, so a saturated budget leaves the
1692			// outbox untouched.
1693			if self.send_in_flight_bytes() >= self.propagation_send_budget() {
1694				// A peer parks once per saturation, not once per tick, so a budget that
1695				// stays full does not grow the deque without bound.
1696				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					// A store read error says nothing about the queued hashes, so the outbox
1718					// is kept and the peer parked: the fetch is retried when a completed send
1719					// or a propagation tick next refills parked peers.
1720					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			// Consume the fetched hashes before the oversized check, otherwise the oversized
1741			// statement would be fetched again on the next iteration.
1742			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				// Everything fetched was filtered out or pruned, but the remaining hashes
1754				// may still yield a chunk.
1755				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	/// Returns whether the result freed the peer's send slot.
1808	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		// A send future is not cancelled on disconnect, so its result can outlive the
1869		// connection. Only the result of the chunk still occupying the slot frees it.
1870		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		// A peer that reconnects inside the send timeout loses its sync on disconnect and gets a
1878		// fresh one under the same `PeerId`; a stale result would advance or abort the wrong sync.
1879		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			// The cursor still points at the unconfirmed chunk, so a later burst resends it.
1885			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		// Reached only for the live sync, which is out of the queue for as long as its chunk is
1896		// in flight; a superseded sync's chunk can still be in flight under the same `PeerId`, so
1897		// the bound is one chunk per sync, not per peer.
1898		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	/// Call when we must propagate ready statements to peers.
1920	async fn propagate_statements(&mut self) {
1921		// Send out statements only when node is not major syncing
1922		if self.sync.is_major_syncing() {
1923			return;
1924		}
1925
1926		// A peer parked by a failed store fetch has no completed send to unpark it, so
1927		// parked peers are also refilled on the tick.
1928		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		// Every entry here belongs to an already drained statement, so it is done
1935		// propagating. Statements imported after the drain get their entries only
1936		// after this clear, when the event loop processes their submit results.
1937		self.recently_received_statements.clear();
1938	}
1939
1940	/// Schedule an initial sync for a peer, sending all known statements.
1941	///
1942	/// This is called both when a new peer connects and when a peer's topic
1943	/// affinity changes (so that newly-matching statements get sent).
1944	/// If the peer already has a pending initial sync, it is replaced.
1945	fn schedule_initial_sync_for_peer(&mut self, peer: PeerId) {
1946		// A peer absent from the map has no entry to mirror the sync's watermark into.
1947		if !self.peers.contains_key(&peer) {
1948			return;
1949		}
1950		// The watermark is read before the existing sync is touched, so a store error
1951		// leaves an in-progress sync running instead of destroying it with no successor.
1952		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			// Hashes queued for propagation before this scheduling sit below the new
1973			// watermark, so the cursor already covers them; dropping the outbox keeps
1974			// them from arriving twice.
1975			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	/// Process pending topic affinity changes for peers that have no active initial sync.
1988	///
1989	/// When a peer sends `ExplicitTopicAffinity`, we defer the expensive
1990	/// `schedule_initial_sync_for_peer` call. This method applies the pending affinity
1991	/// and schedules the sync once the peer's current sync (if any) has completed.
1992	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	/// Set the in-flight initial-sync byte counter.
2012	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	/// Set the in-flight propagation byte counter.
2020	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	/// Total encoded bytes in flight across initial-sync and propagation chunks,
2028	/// held against the shared [`MAX_SEND_IN_FLIGHT_BYTES`] budget.
2029	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	/// Byte budget available to propagation sends.
2035	///
2036	/// While initial syncs are pending, [`INITIAL_SYNC_RESERVED_BYTES`] are withheld:
2037	/// propagation reclaims freed budget synchronously on every completed send, while sync
2038	/// bursts only check on a timer, so without the reserve enough parked propagations
2039	/// starve initial sync indefinitely.
2040	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	/// Refill parked peers' send slots while the in-flight byte budget allows.
2049	///
2050	/// Peers are served in parking order. When the budget saturates the loop stops and the
2051	/// remaining peers keep their position for the next completed send. Each peer gets one
2052	/// attempt per pass: a peer whose store fetch fails parks itself again, and an unbounded
2053	/// loop would spin on it.
2054	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	/// Occupy the peer's send slot with a fresh chunk id and return the id.
2065	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	/// Record initial sync completion metrics for a peer being removed.
2073	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	/// Process one batch of initial sync for the next peer in the queue (round-robin).
2084	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		// A peer whose send slot is busy keeps its turn for a later burst, so one slow peer
2099		// does not stall every other pending sync.
2100		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		// Fetch statements up to max_statement_payload_size, filtering directly in the
2129		// callback (see `fetch_admitted_chunk`).
2130		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				// A store read error says nothing about the journal, and the cursor is
2152				// retained, so the sync resumes from the same position on a later burst.
2153				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		// A failed send must resend the same admissions, so the cursor advances only when a
2163		// send is confirmed; chunks that queue no send advance it here.
2164		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			// Nothing was queued, so no result will arrive for this peer. Put it back and let the
2176			// next burst either send the remainder or observe that the sync is done.
2177			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			// A missing sink usually means the peer is disconnecting, which removes the sync;
2193			// until then the retained cursor lets a later burst retry.
2194			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	/// Default seed used for bloom filters in tests.
2237	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		/// Role to return from `peer_role`. Default: `Full`.
2250		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		/// Admission journal: the vector index is the statement's admission sequence number.
2560		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		/// Insert a statement into the store, recording it in the admission journal.
2575		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		/// Record the hash in the admission journal unless it is already admitted, and
2582		/// return its admission sequence number.
2583		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			// A recent statement is a statement the store holds, so make the drained
2608			// statements visible to `statements_by_hashes` like the real store does.
2609			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				// A journal entry whose statement left the store is a dead sequence number.
2703				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	/// Import one queued statement into the store and feed `result` back into the
2884	/// handler as the main loop would.
2885	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		// Both peers supply the statement while it is queued for validation.
2917		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		// Replaying a propagated statement finds no entry to join, so a peer cannot
2935		// grow the map by resending.
2936		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		// The worker has inserted the statement into the store, but the event loop
2959		// has not processed the validation completion yet.
2960		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		// The second peer sends the same statement inside that window.
2965		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		// Another path imported the statement while it was queued, so validation
3002		// completes with `Known`. The peer that sent it must be recorded all the same.
3003		import_queued_statement(
3004			&mut handler,
3005			&statement_store,
3006			&queue_receiver,
3007			SubmitResult::Known,
3008		)
3009		.await;
3010		// The statement is imported but not yet propagated when another peer forwards it.
3011		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		// The early return skips the drain, so the statement stays recent and its
3046		// entry must stay with it.
3047		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		// 100 KB each, so the tick spans several 1 MiB chunks.
3093		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		// The peer never reads its substream.
3106		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		// Another tick accumulates into the same outbox while the slot is busy.
3114		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		// The pruned statement's hash is queued but the statement left the store.
3143		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(&notification_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		// No new statements arrive for the peer: the tick alone must retry the parked fetch.
3181		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(&notification_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		// The oversized statement heads the outbox. It must be consumed, not
3208		// re-fetched forever, and the statement behind it must still go out.
3209		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(&notification_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		// Several chunks worth of statements with a blocked substream, so the slot
3227		// is taken and a backlog stays queued.
3228		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		// 700 KiB each, so a 1 MiB chunk carries exactly one statement.
3261		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		// Only the first chunk's send fails.
3274		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		// The failure freed the slot and the backlog kept draining: the second
3282		// statement went out, the failed one was not retried.
3283		handler.flush_pending_sends().await;
3284		let sent = get_peer_hashes(&notification_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		// Three fresh statements arrive by tick while the outbox is full and the
3302		// peer's slot is busy.
3303		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		// The hash was appended while the peer's slot was busy, and the peer sent
3353		// us the statement before the slot freed: the encode-time senders check
3354		// must catch what the append-time check could not have seen.
3355		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(&notification_service.get_sent_notifications(), peer_id).is_empty(),
3363			"statement returned to the peer that sent it"
3364		);
3365	}
3366
3367	/// Simulate the network closing the substream for every disconnected
3368	/// peer, so the handler runs its per-peer cleanup.
3369	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			// Manually process statements submission
3418			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),        // Report for first statement
3430				(peer_id, rep::ANY_STATEMENT_REFUND), // Refund for first statement
3431				(peer_id, rep::DUPLICATE_STATEMENT)   // Report for duplicate statement
3432			],
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; // 100KB per statement
3445		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	/// Like `build_handler_no_peers` but the network mock returns `Light` for peer roles.
3603	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		// Create 20MB of statements (200 statements x 100KB each)
3659		// Using 100KB ensures ~10 statements per 1MB batch, requiring ~20 bursts
3660		let num_statements = 200;
3661		let statement_size = 100 * 1024; // 100KB per statement
3662		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			// Use multiple bytes for uniqueness since we have >255 statements
3667			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		// Setup peer and simulate connection
3676		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		// Verify peer was added and initial sync was queued
3688		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		// Process bursts until all statements are sent
3693		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			// Safety limit
3699			assert!(burst_count <= 300, "Too many bursts, possible infinite loop");
3700		}
3701
3702		// Verify multiple bursts were needed
3703		// With 200 statements x 100KB each and ~1MB per batch, we expect many bursts
3704		assert!(
3705			burst_count >= 10,
3706			"Expected multiple bursts for 200 statements of 100KB each, got {}",
3707			burst_count
3708		);
3709
3710		// Verify all statements were sent
3711		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		// Verify cleanup
3733		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		// Two statements small enough to share one chunk, so a statement count cannot be
3744		// mistaken for a chunk count.
3745		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		// The network layer rejects the send with a real error, not a timeout.
3760		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		// The network recovers: the retried chunk delivers the same statements.
3787		notification_service.allow_sends();
3788		handler.process_initial_sync_burst();
3789		handler.flush_pending_sends().await;
3790
3791		let mut sent = get_peer_hashes(&notification_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		// The sync cursor remains the only path, so the statement arrives exactly once.
3820		handler.process_initial_sync_burst();
3821		handler.flush_pending_sends().await;
3822		let sent = get_peer_hashes(&notification_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		// No sink is available for the first burst.
3840		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		// The sink comes back: the retried chunk delivers.
3856		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(&notification_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		// The store recovers: the next burst resumes from the retained cursor.
3892		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(&notification_service.get_sent_notifications(), peer_id);
3897		assert_eq!(sent, vec![hash]);
3898
3899		// The completed send requeued the peer; the next burst observes the finished sync.
3900		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		// 100 KB each, so the batch spans three 1 MiB chunks. Two chunks must remain after the
3911		// sink disappears, otherwise "the whole remainder" and "the current chunk" are the same
3912		// number and the test proves nothing.
3913		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		// Serve the first chunk, then behave as if the peer went away.
3927		notification_service.serve_sinks(1);
3928		handler.propagate_statements().await;
3929		handler.flush_pending_sends().await;
3930
3931		let delivered =
3932			get_peer_hashes(&notification_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		// A topic affinity matching nothing in the store, so the burst finds no
3960		// statement to send.
3961		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		// No chunk means no send result to advance the sync, so the burst has to requeue the peer
3971		// itself or the sync stalls with its entry alive and the peer out of the queue.
3972		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		// 100 KB each, so the store spans several 1 MiB chunks and the peer has more to receive
3989		// after its first one.
3990		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		// The peer never reads its substream.
3999		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		// An affinity change re-schedules the sync while the chunk is still in flight.
4026		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		// Queue a chunk, then arrange for that very send to come back as a failure.
4047		handler.schedule_initial_sync_for_peer(peer_id);
4048		handler.process_initial_sync_burst();
4049		notification_service.fail_sends();
4050
4051		// A reconnect replaces the sync while that chunk is still in flight.
4052		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(&notification_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		// ~1.1 MB of statements, so each peer's first chunk sits just under the 1 MiB cap and a
4113		// handful of peers is enough to exhaust the budget.
4114		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		// No peer reads, so nothing ever leaves the budget.
4123		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		// A throttled burst must leave the round-robin untouched rather than burn a peer's turn.
4143		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		// The whole budget is taken by initial-sync bytes.
4161		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		// A completed initial-sync send frees the budget and refills the parked peer.
4172		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(&notification_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		// One chunk is ~900 KB, so freeing a few bytes admits exactly one of the
4192		// two parked peers into the 16 MiB budget.
4193		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		// Freeing a sliver of budget admits only the first parked peer.
4205		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		// The first chunk's completion frees enough for the second peer.
4217		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		// Propagation holds the whole budget and one peer waits parked.
4243		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		// A completed send frees exactly the reserve: the parked peer must stay parked,
4250		// keeping the headroom for the sync burst.
4251		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		// The sync burst finds the reserved headroom and proceeds.
4263		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(&notification_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		// The next burst observes the drained sync and completes it, releasing the
4275		// reserve back to propagation.
4276		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(&notification_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		// The first queued peer's send slot is taken by a propagation chunk.
4298		handler.in_flight_chunks.insert(peer_ids[0], 42);
4299
4300		handler.process_initial_sync_burst();
4301		handler.flush_pending_sends().await;
4302
4303		// The same burst serves the next queued peer instead of returning.
4304		let sent = notification_service.get_sent_notifications();
4305		assert_eq!(sent.len(), 1);
4306		assert_eq!(sent[0].0, peer_ids[1]);
4307		// The busy peer keeps its sync and its turn.
4308		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		// An initial-sync chunk takes the slot, then a fresh statement arrives by
4319		// tick while the slot is busy.
4320		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		// The initial-sync completion frees the slot and propagation claims it
4336		// before the next burst tick gets a chance.
4337		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(&notification_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		// Propagation bytes alone exhaust the shared budget, so the burst must wait.
4361		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		// Create 20MB of statements (200 statements x 100KB each)
4377		let num_statements = 200;
4378		let statement_size = 100 * 1024; // 100KB per statement
4379		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		// Setup 3 peers and simulate connections
4392		let peer1 = PeerId::random();
4393		let peer2 = PeerId::random();
4394		let peer3 = PeerId::random();
4395
4396		// Connect peers
4397		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		// Verify all peers were added and initial syncs were queued
4409		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		// Track which peer was processed on each burst for round-robin verification
4414		let mut peer_burst_order = Vec::new();
4415		let mut burst_count = 0;
4416
4417		while !handler.pending_initial_syncs.is_empty() {
4418			// Record which peer will be processed next
4419			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			// Safety limit
4426			assert!(burst_count <= 500, "Too many bursts, possible infinite loop");
4427		}
4428
4429		// Verify multiple bursts were needed
4430		// With 3 peers and many bursts per peer, we expect many bursts total
4431		assert!(
4432			burst_count >= 30,
4433			"Expected many bursts for 3 peers with 200 statements each, got {}",
4434			burst_count
4435		);
4436
4437		// Verify round-robin pattern in first 9 bursts (3 peers x 3 rounds)
4438		assert!(peer_burst_order.len() >= 9, "Expected at least 9 bursts");
4439		// First round
4440		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		// Second round
4444		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		// Verify all peers received all statements
4449		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		// Verify cleanup
4464		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		// Calculate the data sizes so that 100 statements together exactly fill max_size.
4474		// This tests that all 100 statements fit in a single notification.
4475		//
4476		// The limit check in `fetch_statement_chunk` is:
4477		//   max_size = MAX_STATEMENT_NOTIFICATION_SIZE - Compact::<u32>::max_encoded_len()
4478		//
4479		// Statement encoding (encodes as Vec<Field>):
4480		// - Compact<u32> for number of fields (1 byte for value 2: expiry + data)
4481		// - Field::Expiry discriminant (1 byte, value 2)
4482		// - u64 expiry value (8 bytes)
4483		// - Field::Data discriminant (1 byte, value 8)
4484		// - Compact<u32> for the data length (2 bytes for small data)
4485		// So per-statement overhead = 1 + 1 + 8 + 1 + 2 = 13 bytes
4486		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; // Vec<Field> length + expiry field + data discriminant + Compact data length
4489		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			// Distribute remainder across first `remainder` statements to exactly fill max_size
4500			let extra = if i < remainder { 1 } else { 0 };
4501			let mut data = vec![42u8; per_statement_data_size + extra];
4502			// Make each statement unique by modifying the first few bytes
4503			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		// Verify our calculation: total encoded size should be <= max_size
4515		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		// All statements should fit in a single chunk
4528		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		// Verify all statements were sent (order may differ due to HashMap iteration)
4553		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		// A filter measuring against the raw `MAX_STATEMENT_NOTIFICATION_SIZE` would admit
4562		// statements the encoder cannot fit, so both must use `max_statement_payload_size`.
4563		let (mut handler, statement_store, _network, notification_service, _, _) = build_handler(0);
4564
4565		// This peer connects as V1 (see negotiated_fallback below).
4566		let payload_limit = max_statement_payload_size(V1_ENVELOPE_OVERHEAD);
4567
4568		// Create first statement that's just over half the payload limit
4569		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		// Create second statement that, combined with the first, exceeds the payload limit.
4575		// This means the filter will only accept the first statement.
4576		let remaining = payload_limit.saturating_sub(stmt1_encoded_size);
4577		let target_stmt2_encoded = remaining + 3; // 3 bytes over limit when combined
4578		let stmt2_data_size = target_stmt2_encoded.saturating_sub(4); // ~4 bytes encoding overhead
4579		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		// Verify our setup: total exceeds payload limit
4586		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		// Setup peer and simulate connection
4599		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		// Verify initial sync was queued covering both statements
4611		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		// Process first burst - should send only one statement (the other doesn't fit)
4616		handler.process_initial_sync_burst();
4617		handler.flush_pending_sends().await;
4618
4619		// The filter and the chunk bound agree, so only the statement that fits is
4620		// fetched and sent.
4621		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		// Verify one of the two statements was sent (order is non-deterministic due to HashMap)
4628		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		// Second statement should still be pending
4635		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		// Process second burst - should send the remaining statement
4640		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		// Both statements should now be sent
4647		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		// The sync is closed by the burst that finds nothing left to send, not by the one that sent
4660		// the last chunk.
4661		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		// Verify peer state was cleaned up
4701		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		// 100 steps of 100ms is 10 simulated seconds: twice the 5s the full burst takes to drain,
4728		// so any drift in the accounting would have tripped by the end.
4729		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	/// Like [`build_handler`], but every peer's rate limiter runs on the returned clock instead of
4834	/// on wall-clock time, so rate-limit behaviour spanning several batches is exact. The quota is
4835	/// the production one.
4836	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	/// The production quota the handler is built with: 50k statements/sec, 250k burst.
4860	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	/// A rate sustained above the quota must drain the bucket across successive batches and
4887	/// eventually get the peer reported and disconnected.
4888	#[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		// No negotiated_fallback means the peer connected on the main protocol (v2).
4945		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		// negotiated_fallback is Some means the peer fell back to v1.
4969		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		// Connect peer as v1 (with fallback).
4995		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		// V1 peer sends raw Vec<Statement>.
5005		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		// Connect peer as v2 (no fallback).
5031		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		// V2 peer sends StatementMessage::Statements.
5041		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		// Connect peer as v2.
5066		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		// Send ExplicitTopicAffinity message.
5081		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		// Affinity is deferred; process it.
5095		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		// The filter should match the topic we inserted.
5103		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		// Connect peer as v2.
5117		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		// Set up topic affinity: peer is interested in topic 0xAA only.
5127		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		// Affinity is deferred; process it.
5141		handler.process_pending_affinities();
5142
5143		// Create statements: one matching, one not matching, one with no topics.
5144		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				// V2 peer gets StatementMessage encoding.
5182				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		// Matching and no-topic statements should be sent; non-matching should be filtered.
5192		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		// Connect peer as v1 (with fallback).
5214		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		// V1 peers have no topic affinity - all statements should be propagated.
5224		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		// Add statements with different topics to the store.
5274		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		// Connect peer as v2.
5296		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		// Light V2 peers should NOT get initial sync on connect (must set affinity first).
5306		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		// Set topic affinity to topic_aa — this triggers the first initial sync.
5312		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		// Affinity is deferred; process it.
5324		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		// Drain initial sync — only stmt_aa and stmt_no_topic should be sent.
5332		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		// Now change affinity to topic_bb — triggers re-sync.
5356		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		// Affinity is deferred; process it.
5368		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		// stmt_bb was previously filtered and should now be sent.
5393		assert!(
5394			sent_hashes_bb.contains(&hash_bb),
5395			"stmt_bb should now be sent after affinity changed to topic_bb"
5396		);
5397		// Known statements are redelivered on affinity change.
5398		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		// This tests the scenario where:
5407		// 1. Peer connects and immediately sets affinity (before initial sync).
5408		// 2. Statements not matching the initial affinity are not delivered.
5409		// 3. When affinity changes to include those topics, they ARE sent.
5410		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		// Also put them in recent_statements so propagate_statements can find them.
5432		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		// Connect peer as v2.
5436		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		// Immediately set affinity to topic_aa BEFORE any initial sync runs.
5446		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		// Affinity is deferred; process it.
5458		handler.process_pending_affinities();
5459
5460		// Drain initial sync — should only send stmt_aa (matches affinity).
5461		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		// Propagation must apply the same affinity filter. The original statements sit below
5484		// the sync watermark and belong to the cursor, so fresh admissions carry the check.
5485		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		// Now change affinity to include topic_bb.
5525		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		// Affinity is deferred; process it.
5540		handler.process_pending_affinities();
5541
5542		// Drain re-sync — stmt_bb should now be sent.
5543		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		// stmt_aa is also redelivered on affinity change.
5564		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		// Full node, V1, no affinity → can receive
5618		assert!(make_peer(false, PeerProtocolVersion::V1, false).can_receive());
5619		// Full node, V2, no affinity → can receive
5620		assert!(make_peer(false, PeerProtocolVersion::V2, false).can_receive());
5621		// Light, V1, no affinity → can receive (V1 doesn't gate)
5622		assert!(make_peer(true, PeerProtocolVersion::V1, false).can_receive());
5623		// Light, V2, no affinity → CANNOT receive (must set affinity first)
5624		assert!(!make_peer(true, PeerProtocolVersion::V2, false).can_receive());
5625		// Light, V2, with affinity → can receive
5626		assert!(make_peer(true, PeerProtocolVersion::V2, true).can_receive());
5627		// Full node, V2, with affinity → can receive
5628		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		// Connect V1 peer.
5640		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		// Connect V2 peer.
5650		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		// Send to V1 peer.
5664		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		// V1 encoding is raw Vec<Statement>.
5674		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		// Send to V2 peer.
5679		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		// V2 encoding is StatementMessage::Statements.
5689		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		// Verify the two encodings are different (V2 has an extra enum discriminant byte).
5697		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		// Add some statements to the store.
5708		let mut stmt1 = new_live_statement();
5709		stmt1.set_plain_data(b"stmt1".to_vec());
5710		statement_store.insert(stmt1);
5711
5712		// Connect peer as V1.
5713		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		// Should have initial sync scheduled.
5723		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		// Add another statement and re-schedule.
5731		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		// Peer should still appear exactly once in the queue (no duplicates).
5738		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		// The new sync should cover both admissions.
5744		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		// Set major syncing to true.
5757		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		// Add a statement so there's something to sync.
5798		let mut stmt = new_live_statement();
5799		stmt.set_plain_data(b"during major sync".to_vec());
5800		statement_store.insert(stmt);
5801
5802		// Add a peer manually.
5803		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		// Scheduling during major sync should queue the peer.
5816		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		// But burst processing should be a no-op while major syncing.
5825		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		// Once major sync completes, burst processing should proceed.
5833		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		// Add statements to the store.
5852		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		// The whole admission journal is covered for redelivery.
5882		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		// The statement is also due for a propagation tick, racing the sync.
5897		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		// A tick between scheduling and the first burst must not queue the statement:
5910		// below the peer's sync watermark it is the cursor's job.
5911		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(&notification_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		// A late tick drains the already synced statement together with a fresh one.
5949		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(&notification_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		// Connect peer as V2.
5974		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		// Send garbage data — should not panic, just log debug.
5984		handler
5985			.handle_notification_event(NotificationEvent::NotificationReceived {
5986				peer: peer_id,
5987				notification: vec![0xFF, 0xFE, 0xFD].into(),
5988			})
5989			.await;
5990
5991		// Send V1-encoded data to V2 peer — also should not panic.
5992		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		// If we got here without panic, the test passes.
6003		assert!(handler.peers.contains_key(&peer_id), "Peer should still be connected");
6004	}
6005
6006	// A batch of one more empty statement (a single `0x00` field-count byte each) than a
6007	// notification can hold. An unbounded `Vec<Statement>` would decode it in full.
6008	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		// `StatementMessage::Statements` variant byte followed by the oversized batch.
6066		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		// V2 has strictly less payload space than V1.
6093		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		// Add a statement so there's something to sync.
6106		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		// Connect as full-node V2 (no fallback, network returns Full role).
6113		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		// Full-node V2 peer should get initial sync immediately (not gated).
6123		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		// Insert 3 statements into recent_statements for propagation
6143		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		// Verify each peer received all 3 statements
6159		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		// Recent statements should be drained
6170		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		// Create 5 statements
6189		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		// peer_a sent s1 and s2, peer_b sent s3, peer_c sent none.
6199		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		// peer_a sent s1,s2 → should only get s3,s4,s5
6219		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		// peer_b sent s3 → should get s1,s2,s4,s5
6227		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		// peer_c sent nothing → should get all 5
6235		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	/// Verifies that peers connecting during major sync are buffered in `deferred_peers` with no
6243	/// network calls, and that a disconnect before sync ends removes the peer from the buffer
6244	#[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		// No network calls while major sync is active
6308		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		// Disconnect before sync ends must remove from buffer only
6313		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		// One remove call must have been issued for the connected peer
6462		{
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		// The recovery peer must be stored and the timeout future must be armed
6473		assert_eq!(handler.sync_recovery_peer, Some(connected_peer));
6474
6475		// Calling try_readd_sync_recovery_peer directly (as the select arm would after the future
6476		// resolves) must re-add the peer and clear the field
6477		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		// Re-entry guard: restore state to simulate a second sync-end while recovery is still
6488		// in flight (sync_recovery_peer is Some). The second call must not issue another remove.
6489		{
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		// flag=false → no recovery
6570		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		// flag=true → recovery fires
6577		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}