referrerpolicy=no-referrer-when-downgrade

cumulus_pallet_parachain_system/
parachain_inherent.rs

1// Copyright (C) Parity Technologies (UK) Ltd.
2// This file is part of Cumulus.
3// SPDX-License-Identifier: Apache-2.0
4
5// Licensed under the Apache License, Version 2.0 (the "License");
6// you may not use this file except in compliance with the License.
7// You may obtain a copy of the License at
8//
9// 	http://www.apache.org/licenses/LICENSE-2.0
10//
11// Unless required by applicable law or agreed to in writing, software
12// distributed under the License is distributed on an "AS IS" BASIS,
13// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
14// See the License for the specific language governing permissions and
15// limitations under the License.
16
17//! Cumulus parachain inherent related structures.
18
19use alloc::{collections::btree_map::BTreeMap, vec, vec::Vec};
20use core::fmt::Debug;
21use cumulus_primitives_core::{
22	relay_chain::{
23		ApprovedPeerId, BlockNumber as RelayChainBlockNumber, BlockNumber, Header as RelayHeader,
24	},
25	InboundDownwardMessage, InboundHrmpMessage, ParaId, PersistedValidationData,
26};
27use cumulus_primitives_parachain_inherent::{HashedMessage, ParachainInherentData};
28use frame_support::{
29	defensive,
30	pallet_prelude::{Decode, DecodeWithMemTracking, Encode},
31};
32use scale_info::TypeInfo;
33use sp_core::{bounded::BoundedSlice, Get};
34
35/// A structure that helps identify a message inside a collection of messages sorted by `sent_at`.
36///
37/// This structure contains a `sent_at` field and a reverse index. Using this information, we can
38/// identify a message inside a sorted collection by walking back `reverse_idx` positions starting
39/// from the last message that has the provided `sent_at`.
40///
41/// We use a reverse index instead of a normal index because sometimes the messages at the
42/// beginning of the collection are being pruned.
43///
44/// # Example
45///
46///
47/// For the collection
48/// `msgs = [{sent_at: 0}, {sent_at: 1}, {sent_at: 1}, {sent_at: 1}, {sent_at: 1}, {sent_at: 3}]`
49///
50/// `InboundMessageId {sent_at: 1, reverse_idx: 0}` points to `msgs[4]`
51/// `InboundMessageId {sent_at: 1, reverse_idx: 3}` points to `msgs[1]`
52/// `InboundMessageId {sent_at: 1, reverse_idx: 4}` points to `msgs[0]`
53#[derive(Encode, Decode, DecodeWithMemTracking, Clone, Default, Debug, PartialEq, TypeInfo)]
54pub struct InboundMessageId {
55	/// The block number at which this message was added to the message passing queue
56	/// on the relay chain.
57	pub sent_at: BlockNumber,
58	/// The reverse index of the message in the collection of messages sent at `sent_at`.
59	pub reverse_idx: u32,
60}
61
62/// A message that was received by the parachain.
63pub trait InboundMessage {
64	/// The corresponding compressed message.
65	/// This should be an equivalent message that stores the same metadata as the current message,
66	/// but stores only a hash of the message data.
67	type CompressedMessage: Debug;
68
69	/// Gets the message data.
70	fn data(&self) -> &[u8];
71
72	/// Gets the relay chain number where the current message was pushed to the corresponding
73	/// relay chain queue.
74	fn sent_at(&self) -> RelayChainBlockNumber;
75
76	/// Converts the current message into a `CompressedMessage`
77	fn to_compressed(&self) -> Self::CompressedMessage;
78}
79
80/// A collection of inbound messages.
81#[derive(
82	codec::Encode, codec::Decode, codec::DecodeWithMemTracking, Debug, Clone, PartialEq, TypeInfo,
83)]
84pub struct InboundMessagesCollection<Message: InboundMessage> {
85	messages: Vec<Message>,
86}
87
88impl<Message: InboundMessage> InboundMessagesCollection<Message> {
89	/// Creates a new instance of `InboundMessagesCollection` that contains the provided `messages`.
90	pub fn new(messages: Vec<Message>) -> Self {
91		Self { messages }
92	}
93
94	/// Drop all the messages up to `last_processed_msg`.
95	pub fn drop_processed_messages(&mut self, last_processed_msg: &InboundMessageId) {
96		let mut last_processed_msg_idx = None;
97		let messages = &mut self.messages;
98		for (idx, message) in messages.iter().enumerate().rev() {
99			let sent_at = message.sent_at();
100			if sent_at == last_processed_msg.sent_at {
101				last_processed_msg_idx = idx.checked_sub(last_processed_msg.reverse_idx as usize);
102				break;
103			}
104			// If we build on the same relay parent twice, we will receive the same messages again
105			// while `last_processed_msg` may have been increased. We need this check to make sure
106			// that the old messages are dropped.
107			if sent_at < last_processed_msg.sent_at {
108				last_processed_msg_idx = Some(idx);
109				break;
110			}
111		}
112		if let Some(last_processed_msg_idx) = last_processed_msg_idx {
113			messages.drain(..=last_processed_msg_idx);
114		}
115	}
116
117	/// Converts `self` into an [`AbridgedInboundMessagesCollection`].
118	///
119	/// The first messages in `self` (up to the provided `size_limit`) are kept in their current
120	/// form (they will contain the full message data).
121	/// The messages that exceed that limit are hashed.
122	pub fn into_abridged(
123		self,
124		size_limit: &mut usize,
125	) -> AbridgedInboundMessagesCollection<Message> {
126		let mut messages = self.messages;
127
128		let mut split_off_pos = messages.len();
129		for (idx, message) in messages.iter().enumerate() {
130			if *size_limit < message.data().len() {
131				break;
132			}
133			*size_limit -= message.data().len();
134
135			split_off_pos = idx + 1;
136		}
137
138		let extra_messages = messages.split_off(split_off_pos);
139		let hashed_messages = extra_messages.iter().map(|msg| msg.to_compressed()).collect();
140
141		AbridgedInboundMessagesCollection { full_messages: messages, hashed_messages }
142	}
143}
144
145/// A struct containing some info about the expected size of the abridged inbound messages.
146pub struct AbridgedInboundMessagesSizeInfo {
147	/// The max size of the full messages collection
148	pub max_full_messages_size: usize,
149	/// The max size of the first hashed message
150	pub first_hashed_msg_max_size: usize,
151}
152
153/// A compressed collection of inbound messages.
154///
155/// The first messages in the collection (up to a limit) contain the full message data.
156/// The messages that exceed that limit are hashed.
157#[derive(
158	codec::Encode, codec::Decode, codec::DecodeWithMemTracking, Debug, Clone, PartialEq, TypeInfo,
159)]
160pub struct AbridgedInboundMessagesCollection<Message: InboundMessage> {
161	full_messages: Vec<Message>,
162	hashed_messages: Vec<Message::CompressedMessage>,
163}
164
165impl<Message: InboundMessage> AbridgedInboundMessagesCollection<Message> {
166	/// Gets a tuple containing both the full messages and the hashed messages
167	/// stored by the current collection.
168	pub fn messages(&self) -> (&[Message], &[Message::CompressedMessage]) {
169		(&self.full_messages, &self.hashed_messages)
170	}
171
172	/// Check that the current collection contains at least 1 full message if needed.
173	pub fn check_enough_messages_included_basic(&self, collection_name: &str) {
174		if self.hashed_messages.is_empty() {
175			return;
176		}
177
178		// Here we just check that there is at least 1 full message.
179		assert!(
180			self.full_messages.len() >= 1,
181			"[{}] Advancement rule violation: full messages missing",
182			collection_name,
183		);
184	}
185
186	/// Check that the current collection contains as many full messages as possible, taking into
187	/// consideration the collection constraints.
188	///
189	/// The `AbridgedInboundMessagesCollection` is provided to the runtime by a collator.
190	/// A malicious collator can provide a collection that contains no full messages or fewer
191	/// full messages than possible, leading to censorship.
192	pub fn check_enough_messages_included_advanced(
193		&self,
194		collection_name: &str,
195		size_info: AbridgedInboundMessagesSizeInfo,
196	) {
197		// We should check that the collection contains as many full messages as possible
198		// without exceeding the max expected size.
199		let AbridgedInboundMessagesSizeInfo { max_full_messages_size, first_hashed_msg_max_size } =
200			size_info;
201
202		let mut full_messages_size = 0usize;
203		for msg in &self.full_messages {
204			full_messages_size = full_messages_size.saturating_add(msg.data().len());
205		}
206
207		// The worst case scenario is that were the first message that had to be hashed
208		// is a max size message.
209		assert!(
210			full_messages_size.saturating_add(first_hashed_msg_max_size) > max_full_messages_size,
211			"[{}] Advancement rule violation: full messages size smaller than expected. \
212			full msgs size: {}, first hashed msg max size: {}, max full msgs size: {}",
213			collection_name,
214			full_messages_size,
215			first_hashed_msg_max_size,
216			max_full_messages_size
217		);
218	}
219}
220
221impl<Message: InboundMessage> Default for AbridgedInboundMessagesCollection<Message> {
222	fn default() -> Self {
223		Self { full_messages: vec![], hashed_messages: vec![] }
224	}
225}
226
227impl InboundMessage for InboundDownwardMessage<RelayChainBlockNumber> {
228	type CompressedMessage = HashedMessage;
229
230	fn data(&self) -> &[u8] {
231		&self.msg
232	}
233
234	fn sent_at(&self) -> RelayChainBlockNumber {
235		self.sent_at
236	}
237
238	fn to_compressed(&self) -> Self::CompressedMessage {
239		self.into()
240	}
241}
242
243pub type InboundDownwardMessages =
244	InboundMessagesCollection<InboundDownwardMessage<RelayChainBlockNumber>>;
245
246pub type AbridgedInboundDownwardMessages =
247	AbridgedInboundMessagesCollection<InboundDownwardMessage<RelayChainBlockNumber>>;
248
249impl AbridgedInboundDownwardMessages {
250	/// Returns an iterator over the messages that maps them to `BoundedSlices`.
251	pub fn bounded_msgs_iter<MaxMessageLen: Get<u32>>(
252		&self,
253	) -> impl Iterator<Item = BoundedSlice<'_, u8, MaxMessageLen>> {
254		self.full_messages
255			.iter()
256			// Note: we are not using `.defensive()` here since that prints the whole value to
257			// console. In case that the message is too long, this clogs up the log quite badly.
258			.filter_map(|m| match BoundedSlice::try_from(&m.msg[..]) {
259				Ok(bounded) => Some(bounded),
260				Err(_) => {
261					defensive!("Inbound Downward message was too long; dropping");
262					None
263				},
264			})
265	}
266}
267
268impl InboundMessage for (ParaId, InboundHrmpMessage) {
269	type CompressedMessage = (ParaId, HashedMessage);
270
271	fn data(&self) -> &[u8] {
272		&self.1.data
273	}
274
275	fn sent_at(&self) -> RelayChainBlockNumber {
276		self.1.sent_at
277	}
278
279	fn to_compressed(&self) -> Self::CompressedMessage {
280		let (sender, message) = self;
281		(*sender, message.into())
282	}
283}
284
285/// Similar to [`InboundMessageId`], but also containing the sending parachain id.
286#[derive(Encode, Decode, DecodeWithMemTracking, Clone, Debug, PartialEq, TypeInfo)]
287pub enum InboundHrmpMessageId {
288	Generic(InboundMessageId),
289	Specific {
290		/// The block number at which this message was added to the message passing queue
291		/// on the relay chain.
292		sent_at: BlockNumber,
293		/// The sending parachain id.
294		sender: ParaId,
295		/// The reverse index of the message in the collection of messages sent at `sent_at`
296		/// by `sender`.
297		reverse_idx: u32,
298	},
299}
300
301impl InboundHrmpMessageId {
302	pub fn sent_at(&self) -> BlockNumber {
303		match self {
304			InboundHrmpMessageId::Generic(id) => id.sent_at,
305			InboundHrmpMessageId::Specific { sent_at, .. } => *sent_at,
306		}
307	}
308
309	pub fn sender(&self) -> Option<ParaId> {
310		match self {
311			InboundHrmpMessageId::Generic(_) => None,
312			InboundHrmpMessageId::Specific { sender, .. } => Some(*sender),
313		}
314	}
315
316	pub fn inc_reverse_idx(&mut self) {
317		let reverse_idx = match self {
318			InboundHrmpMessageId::Generic(id) => &mut id.reverse_idx,
319			InboundHrmpMessageId::Specific { reverse_idx, .. } => reverse_idx,
320		};
321
322		*reverse_idx += 1;
323	}
324}
325
326pub type InboundHrmpMessages = InboundMessagesCollection<(ParaId, InboundHrmpMessage)>;
327
328impl InboundHrmpMessages {
329	// Prepare horizontal messages for a more convenient processing:
330	//
331	// Instead of a mapping from a para to a list of inbound HRMP messages, we will have a
332	// list of tuples `(sender, message)` first ordered by `sent_at` (the relay chain block
333	// number in which the message hit the relay-chain) and second ordered by para id
334	// ascending.
335	pub fn from_map(messages_map: BTreeMap<ParaId, Vec<InboundHrmpMessage>>) -> Self {
336		let mut messages = messages_map
337			.into_iter()
338			.flat_map(|(sender, channel_contents)| {
339				channel_contents.into_iter().map(move |message| (sender, message))
340			})
341			.collect::<Vec<_>>();
342		messages.sort_by(|(sender_a, msg_a), (sender_b, msg_b)| {
343			// first sort by sent-at and then by the para id
344			(msg_a.sent_at, sender_a).cmp(&(msg_b.sent_at, sender_b))
345		});
346
347		Self { messages }
348	}
349
350	/// Drop all the messages up to `last_processed_msg`.
351	pub fn drop_hrmp_processed_messages(&mut self, last_processed_msg: &InboundHrmpMessageId) {
352		let (input_sent_at, input_sender, input_reverse_idx) = match last_processed_msg {
353			InboundHrmpMessageId::Generic(id) => {
354				return self.drop_processed_messages(id);
355			},
356			InboundHrmpMessageId::Specific { sent_at, sender, reverse_idx } => {
357				(sent_at, sender, reverse_idx)
358			},
359		};
360
361		let mut last_processed_msg_idx = None;
362		let messages = &mut self.messages;
363		for (idx, (sender, message)) in messages.iter().enumerate().rev() {
364			let sent_at = message.sent_at;
365			if sent_at == *input_sent_at && input_sender == sender {
366				last_processed_msg_idx = idx.checked_sub(*input_reverse_idx as usize);
367				break;
368			}
369			// If we build on the same relay parent twice, we will receive the same messages again
370			// while `last_processed_msg` may have been increased.
371			// Also, if an HRMP channel was closed, the messages from that sender will not be sent,
372			// even if our `last_processed_msg` point to a message from it.
373			if sent_at < *input_sent_at || (sender < input_sender && sent_at == *input_sent_at) {
374				last_processed_msg_idx = Some(idx);
375				break;
376			}
377		}
378		if let Some(last_processed_msg_idx) = last_processed_msg_idx {
379			messages.drain(..=last_processed_msg_idx);
380		}
381	}
382}
383
384pub type AbridgedInboundHrmpMessages =
385	AbridgedInboundMessagesCollection<(ParaId, InboundHrmpMessage)>;
386
387impl AbridgedInboundHrmpMessages {
388	/// Returns an iterator over the deconstructed messages.
389	pub fn flat_msgs_iter(&self) -> impl Iterator<Item = (ParaId, RelayChainBlockNumber, &[u8])> {
390		self.full_messages
391			.iter()
392			.map(|&(sender, ref message)| (sender, message.sent_at, &message.data[..]))
393	}
394}
395
396/// The basic inherent data that is passed by the collator to the parachain runtime.
397/// This data doesn't contain any messages.
398#[derive(
399	codec::Encode, codec::Decode, codec::DecodeWithMemTracking, Debug, Clone, PartialEq, TypeInfo,
400)]
401pub struct BasicParachainInherentData {
402	pub validation_data: PersistedValidationData,
403	pub relay_chain_state: sp_trie::StorageProof,
404	pub relay_parent_descendants: Vec<RelayHeader>,
405	pub collator_peer_id: Option<ApprovedPeerId>,
406}
407
408/// The messages that are passed by the collator to the parachain runtime as part of the
409/// inherent data.
410#[derive(
411	codec::Encode, codec::Decode, codec::DecodeWithMemTracking, Debug, Clone, PartialEq, TypeInfo,
412)]
413pub struct InboundMessagesData {
414	pub downward_messages: AbridgedInboundDownwardMessages,
415	pub horizontal_messages: AbridgedInboundHrmpMessages,
416}
417
418impl InboundMessagesData {
419	/// Creates a new instance of `InboundMessagesData` with the provided messages.
420	pub fn new(
421		dmq_msgs: AbridgedInboundDownwardMessages,
422		hrmp_msgs: AbridgedInboundHrmpMessages,
423	) -> Self {
424		Self { downward_messages: dmq_msgs, horizontal_messages: hrmp_msgs }
425	}
426}
427
428/// Deconstructs a `ParachainInherentData` instance.
429pub fn deconstruct_parachain_inherent_data(
430	data: ParachainInherentData,
431) -> (BasicParachainInherentData, InboundDownwardMessages, InboundHrmpMessages) {
432	(
433		BasicParachainInherentData {
434			validation_data: data.validation_data,
435			relay_chain_state: data.relay_chain_state,
436			relay_parent_descendants: data.relay_parent_descendants,
437			collator_peer_id: data.collator_peer_id,
438		},
439		InboundDownwardMessages::new(data.downward_messages),
440		InboundHrmpMessages::from_map(data.horizontal_messages),
441	)
442}
443
444#[cfg(test)]
445mod tests {
446	use super::*;
447
448	fn build_inbound_dm_vec(
449		info: &[(RelayChainBlockNumber, usize)],
450	) -> Vec<InboundDownwardMessage<RelayChainBlockNumber>> {
451		let mut messages = vec![];
452		for (sent_at, size) in info.iter() {
453			let data = vec![1; *size];
454			messages.push(InboundDownwardMessage { sent_at: *sent_at, msg: data })
455		}
456		messages
457	}
458
459	#[test]
460	fn drop_hrmp_processed_messages_works() {
461		let msgs_vec = vec![
462			// sent_at: 0
463			(0.into(), InboundHrmpMessage { sent_at: 0, data: vec![1] }),
464			(2.into(), InboundHrmpMessage { sent_at: 0, data: vec![2] }),
465			(2.into(), InboundHrmpMessage { sent_at: 0, data: vec![3] }),
466			// sent_at: 1
467			(0.into(), InboundHrmpMessage { sent_at: 1, data: vec![4] }),
468			(2.into(), InboundHrmpMessage { sent_at: 1, data: vec![5] }),
469			(2.into(), InboundHrmpMessage { sent_at: 1, data: vec![6] }),
470		];
471
472		// until sent_at: 0
473
474		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
475		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
476			sent_at: 0,
477			sender: 0.into(),
478			reverse_idx: 2,
479		});
480		assert_eq!(msgs.messages, msgs_vec[..]);
481
482		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
483		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
484			sent_at: 0,
485			sender: 0.into(),
486			reverse_idx: 1,
487		});
488		assert_eq!(msgs.messages, msgs_vec[..]);
489
490		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
491		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
492			sent_at: 0,
493			sender: 0.into(),
494			reverse_idx: 0,
495		});
496		assert_eq!(msgs.messages, msgs_vec[1..]);
497
498		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
499		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
500			sent_at: 0,
501			sender: 1.into(),
502			reverse_idx: 2,
503		});
504		assert_eq!(msgs.messages, msgs_vec[1..]);
505
506		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
507		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
508			sent_at: 0,
509			sender: 2.into(),
510			reverse_idx: 2,
511		});
512		assert_eq!(msgs.messages, msgs_vec[1..]);
513
514		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
515		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
516			sent_at: 0,
517			sender: 2.into(),
518			reverse_idx: 1,
519		});
520		assert_eq!(msgs.messages, msgs_vec[2..]);
521
522		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
523		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
524			sent_at: 0,
525			sender: 2.into(),
526			reverse_idx: 0,
527		});
528		assert_eq!(msgs.messages, msgs_vec[3..]);
529
530		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
531		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
532			sent_at: 0,
533			sender: 3.into(),
534			reverse_idx: 1,
535		});
536		assert_eq!(msgs.messages, msgs_vec[3..]);
537
538		// until sent_at: 1
539
540		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
541		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
542			sent_at: 1,
543			sender: 0.into(),
544			reverse_idx: 1,
545		});
546		assert_eq!(msgs.messages, msgs_vec[3..]);
547
548		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
549		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
550			sent_at: 1,
551			sender: 0.into(),
552			reverse_idx: 0,
553		});
554		assert_eq!(msgs.messages, msgs_vec[4..]);
555
556		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
557		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
558			sent_at: 1,
559			sender: 1.into(),
560			reverse_idx: 2,
561		});
562		assert_eq!(msgs.messages, msgs_vec[4..]);
563
564		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
565		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
566			sent_at: 1,
567			sender: 2.into(),
568			reverse_idx: 2,
569		});
570		assert_eq!(msgs.messages, msgs_vec[4..]);
571
572		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
573		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
574			sent_at: 1,
575			sender: 2.into(),
576			reverse_idx: 1,
577		});
578		assert_eq!(msgs.messages, msgs_vec[5..]);
579
580		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
581		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
582			sent_at: 1,
583			sender: 2.into(),
584			reverse_idx: 0,
585		});
586		assert_eq!(msgs.messages, msgs_vec[6..]);
587
588		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
589		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Specific {
590			sent_at: 1,
591			sender: 3.into(),
592			reverse_idx: 1,
593		});
594		assert_eq!(msgs.messages, msgs_vec[6..]);
595	}
596
597	#[test]
598	fn drop_processed_messages_works() {
599		let msgs_vec = vec![
600			(0.into(), InboundHrmpMessage { sent_at: 0, data: vec![1] }),
601			(0.into(), InboundHrmpMessage { sent_at: 0, data: vec![2] }),
602			(0.into(), InboundHrmpMessage { sent_at: 2, data: vec![3] }),
603			(0.into(), InboundHrmpMessage { sent_at: 2, data: vec![4] }),
604			(0.into(), InboundHrmpMessage { sent_at: 2, data: vec![5] }),
605			(0.into(), InboundHrmpMessage { sent_at: 2, data: vec![6] }),
606			(0.into(), InboundHrmpMessage { sent_at: 3, data: vec![7] }),
607		];
608
609		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
610		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Generic(InboundMessageId {
611			sent_at: 3,
612			reverse_idx: 0,
613		}));
614		assert_eq!(msgs.messages, []);
615
616		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
617		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Generic(InboundMessageId {
618			sent_at: 2,
619			reverse_idx: 0,
620		}));
621		assert_eq!(msgs.messages, msgs_vec[6..]);
622
623		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
624		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Generic(InboundMessageId {
625			sent_at: 2,
626			reverse_idx: 1,
627		}));
628		assert_eq!(msgs.messages, msgs_vec[5..]);
629
630		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
631		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Generic(InboundMessageId {
632			sent_at: 2,
633			reverse_idx: 4,
634		}));
635		assert_eq!(msgs.messages, msgs_vec[2..]);
636
637		// Go back starting from the last message sent at block 2, with 1 more message than the
638		// total number of messages sent at 2.
639		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
640		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Generic(InboundMessageId {
641			sent_at: 2,
642			reverse_idx: 5,
643		}));
644		assert_eq!(msgs.messages, msgs_vec[1..]);
645
646		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
647		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Generic(InboundMessageId {
648			sent_at: 0,
649			reverse_idx: 1,
650		}));
651		assert_eq!(msgs.messages, msgs_vec[1..]);
652
653		// Go back starting from the last message sent at block 0, with 1 more message than the
654		// total number of messages sent at 0.
655		let mut msgs = InboundHrmpMessages::new(msgs_vec.clone());
656		msgs.drop_hrmp_processed_messages(&InboundHrmpMessageId::Generic(InboundMessageId {
657			sent_at: 0,
658			reverse_idx: 3,
659		}));
660		assert_eq!(msgs.messages, msgs_vec);
661	}
662
663	#[test]
664	fn into_abridged_works() {
665		let msgs = InboundDownwardMessages::new(vec![]);
666		let mut size_limit = 0;
667		let abridged_msgs = msgs.into_abridged(&mut size_limit);
668		assert_eq!(size_limit, 0);
669		assert_eq!(&abridged_msgs.full_messages, &vec![]);
670		assert_eq!(abridged_msgs.hashed_messages, vec![]);
671
672		let msgs_vec = build_inbound_dm_vec(&[(0, 100), (0, 100), (0, 150), (0, 50)]);
673		let msgs = InboundDownwardMessages::new(msgs_vec.clone());
674
675		let mut size_limit = 150;
676		let abridged_msgs = msgs.clone().into_abridged(&mut size_limit);
677		assert_eq!(size_limit, 50);
678		assert_eq!(&abridged_msgs.full_messages, &msgs_vec[..1]);
679		assert_eq!(
680			abridged_msgs.hashed_messages,
681			vec![(&msgs_vec[1]).into(), (&msgs_vec[2]).into(), (&msgs_vec[3]).into()]
682		);
683
684		let mut size_limit = 200;
685		let abridged_msgs = msgs.clone().into_abridged(&mut size_limit);
686		assert_eq!(size_limit, 0);
687		assert_eq!(&abridged_msgs.full_messages, &msgs_vec[..2]);
688		assert_eq!(
689			abridged_msgs.hashed_messages,
690			vec![(&msgs_vec[2]).into(), (&msgs_vec[3]).into()]
691		);
692
693		let mut size_limit = 399;
694		let abridged_msgs = msgs.clone().into_abridged(&mut size_limit);
695		assert_eq!(size_limit, 49);
696		assert_eq!(&abridged_msgs.full_messages, &msgs_vec[..3]);
697		assert_eq!(abridged_msgs.hashed_messages, vec![(&msgs_vec[3]).into()]);
698
699		let mut size_limit = 400;
700		let abridged_msgs = msgs.clone().into_abridged(&mut size_limit);
701		assert_eq!(size_limit, 0);
702		assert_eq!(&abridged_msgs.full_messages, &msgs_vec);
703		assert_eq!(abridged_msgs.hashed_messages, vec![]);
704	}
705
706	#[test]
707	fn from_map_works() {
708		let mut messages_map: BTreeMap<ParaId, Vec<InboundHrmpMessage>> = BTreeMap::new();
709		messages_map.insert(
710			1000.into(),
711			vec![
712				InboundHrmpMessage { sent_at: 0, data: vec![0] },
713				InboundHrmpMessage { sent_at: 0, data: vec![1] },
714				InboundHrmpMessage { sent_at: 1, data: vec![2] },
715			],
716		);
717		messages_map.insert(
718			2000.into(),
719			vec![
720				InboundHrmpMessage { sent_at: 0, data: vec![3] },
721				InboundHrmpMessage { sent_at: 0, data: vec![4] },
722				InboundHrmpMessage { sent_at: 1, data: vec![5] },
723			],
724		);
725		messages_map.insert(
726			3000.into(),
727			vec![
728				InboundHrmpMessage { sent_at: 0, data: vec![6] },
729				InboundHrmpMessage { sent_at: 1, data: vec![7] },
730				InboundHrmpMessage { sent_at: 2, data: vec![8] },
731				InboundHrmpMessage { sent_at: 3, data: vec![9] },
732				InboundHrmpMessage { sent_at: 4, data: vec![10] },
733			],
734		);
735
736		let msgs = InboundHrmpMessages::from_map(messages_map);
737		assert_eq!(
738			msgs.messages,
739			[
740				(1000.into(), InboundHrmpMessage { sent_at: 0, data: vec![0] }),
741				(1000.into(), InboundHrmpMessage { sent_at: 0, data: vec![1] }),
742				(2000.into(), InboundHrmpMessage { sent_at: 0, data: vec![3] }),
743				(2000.into(), InboundHrmpMessage { sent_at: 0, data: vec![4] }),
744				(3000.into(), InboundHrmpMessage { sent_at: 0, data: vec![6] }),
745				(1000.into(), InboundHrmpMessage { sent_at: 1, data: vec![2] }),
746				(2000.into(), InboundHrmpMessage { sent_at: 1, data: vec![5] }),
747				(3000.into(), InboundHrmpMessage { sent_at: 1, data: vec![7] }),
748				(3000.into(), InboundHrmpMessage { sent_at: 2, data: vec![8] }),
749				(3000.into(), InboundHrmpMessage { sent_at: 3, data: vec![9] }),
750				(3000.into(), InboundHrmpMessage { sent_at: 4, data: vec![10] })
751			]
752		)
753	}
754
755	#[test]
756	fn check_enough_messages_included_basic_works() {
757		let mut messages = AbridgedInboundHrmpMessages {
758			full_messages: vec![(
759				1000.into(),
760				InboundHrmpMessage { sent_at: 0, data: vec![1; 100] },
761			)],
762			hashed_messages: vec![(
763				2000.into(),
764				HashedMessage { sent_at: 1, msg_hash: Default::default() },
765			)],
766		};
767
768		messages.check_enough_messages_included_basic("Test");
769
770		messages.full_messages = vec![];
771		let result =
772			std::panic::catch_unwind(|| messages.check_enough_messages_included_basic("Test"));
773		assert!(result.is_err());
774
775		messages.hashed_messages = vec![];
776		messages.check_enough_messages_included_basic("Test");
777	}
778
779	#[test]
780	fn check_enough_messages_included_advanced_works() {
781		let mixed_messages = AbridgedInboundHrmpMessages {
782			full_messages: vec![(
783				1000.into(),
784				InboundHrmpMessage { sent_at: 0, data: vec![1; 50] },
785			)],
786			hashed_messages: vec![(
787				2000.into(),
788				HashedMessage { sent_at: 1, msg_hash: Default::default() },
789			)],
790		};
791		let result = std::panic::catch_unwind(|| {
792			mixed_messages.check_enough_messages_included_advanced(
793				"Test",
794				AbridgedInboundMessagesSizeInfo {
795					max_full_messages_size: 100,
796					first_hashed_msg_max_size: 50,
797				},
798			)
799		});
800		assert!(result.is_err());
801		mixed_messages.check_enough_messages_included_advanced(
802			"Test",
803			AbridgedInboundMessagesSizeInfo {
804				max_full_messages_size: 100,
805				first_hashed_msg_max_size: 51,
806			},
807		);
808	}
809}