cumulus_client_consensus_common/
lib.rs1use codec::Decode;
19use polkadot_primitives::{Block as PBlock, Hash as PHash, Header as PHeader, ValidationCodeHash};
20
21use cumulus_primitives_core::{relay_chain, AbridgedHostConfiguration};
22use cumulus_relay_chain_interface::{RelayChainError, RelayChainInterface};
23
24use sc_client_api::Backend;
25use sc_consensus::{shared_data::SharedData, BlockImport, ImportResult};
26use sp_consensus_slots::Slot;
27
28use sp_runtime::traits::{Block as BlockT, Header as HeaderT};
29use sp_timestamp::Timestamp;
30
31use std::{sync::Arc, time::Duration};
32
33mod finality;
34mod level_monitor;
35mod parachain_consensus;
36mod parent_search;
37#[cfg(test)]
38mod tests;
39
40pub use finality::old_finalized_hash;
41pub use parent_search::*;
42
43pub use cumulus_relay_chain_streams::finalized_heads;
44pub use parachain_consensus::spawn_parachain_consensus_tasks;
45
46use level_monitor::LevelMonitor;
47pub use level_monitor::{LevelLimit, MAX_LEAVES_PER_LEVEL_SENSIBLE_DEFAULT};
48
49pub mod import_queue;
50
51const LOG_TARGET: &str = "consensus::common";
52
53pub trait ValidationCodeHashProvider<Hash> {
56 fn code_hash_at(&self, at: Hash) -> Option<ValidationCodeHash>;
57}
58
59impl<F, Hash> ValidationCodeHashProvider<Hash> for F
60where
61 F: Fn(Hash) -> Option<ValidationCodeHash>,
62{
63 fn code_hash_at(&self, at: Hash) -> Option<ValidationCodeHash> {
64 (self)(at)
65 }
66}
67
68pub struct ParachainCandidate<B> {
70 pub block: B,
72 pub proof: sp_trie::StorageProof,
74}
75
76pub struct ParachainBlockImport<Block: BlockT, BI, BE> {
84 inner: BI,
85 monitor: Option<SharedData<LevelMonitor<Block, BE>>>,
86 delayed_best_block: bool,
87}
88
89impl<Block: BlockT, BI, BE: Backend<Block>> ParachainBlockImport<Block, BI, BE> {
90 pub fn new(inner: BI, backend: Arc<BE>) -> Self {
94 Self::new_with_limit(inner, backend, LevelLimit::Default)
95 }
96
97 pub fn new_with_limit(inner: BI, backend: Arc<BE>, level_leaves_max: LevelLimit) -> Self {
102 let level_limit = match level_leaves_max {
103 LevelLimit::None => None,
104 LevelLimit::Some(limit) => Some(limit),
105 LevelLimit::Default => Some(MAX_LEAVES_PER_LEVEL_SENSIBLE_DEFAULT),
106 };
107
108 let monitor =
109 level_limit.map(|level_limit| SharedData::new(LevelMonitor::new(level_limit, backend)));
110
111 Self { inner, monitor, delayed_best_block: false }
112 }
113
114 pub fn new_with_delayed_best_block(inner: BI, backend: Arc<BE>) -> Self {
118 Self {
119 delayed_best_block: true,
120 ..Self::new_with_limit(inner, backend, LevelLimit::Default)
121 }
122 }
123}
124
125impl<Block: BlockT, I: Clone, BE> Clone for ParachainBlockImport<Block, I, BE> {
126 fn clone(&self) -> Self {
127 ParachainBlockImport {
128 inner: self.inner.clone(),
129 monitor: self.monitor.clone(),
130 delayed_best_block: self.delayed_best_block,
131 }
132 }
133}
134
135#[async_trait::async_trait]
136impl<Block, BI, BE> BlockImport<Block> for ParachainBlockImport<Block, BI, BE>
137where
138 Block: BlockT,
139 BI: BlockImport<Block> + Send + Sync,
140 BE: Backend<Block>,
141{
142 type Error = BI::Error;
143
144 async fn check_block(
145 &self,
146 block: sc_consensus::BlockCheckParams<Block>,
147 ) -> Result<sc_consensus::ImportResult, Self::Error> {
148 self.inner.check_block(block).await
149 }
150
151 async fn import_block(
152 &self,
153 mut params: sc_consensus::BlockImportParams<Block>,
154 ) -> Result<sc_consensus::ImportResult, Self::Error> {
155 let hash = params.post_hash();
157 let number = *params.header.number();
158
159 if params.with_state() {
160 params.finalized = true;
164 }
165
166 if self.delayed_best_block {
167 params.fork_choice = Some(sc_consensus::ForkChoiceStrategy::Custom(
170 params.origin == sp_consensus::BlockOrigin::NetworkInitialSync,
171 ));
172 }
173
174 let maybe_lock = self.monitor.as_ref().map(|monitor_lock| {
175 let mut monitor = monitor_lock.shared_data_locked();
176 monitor.enforce_limit(number);
177 monitor.release_mutex()
178 });
179
180 let res = self.inner.import_block(params).await?;
181
182 if let (Some(mut monitor_lock), ImportResult::Imported(_)) = (maybe_lock, &res) {
183 let mut monitor = monitor_lock.upgrade();
184 monitor.block_imported(number, hash);
185 }
186
187 Ok(res)
188 }
189}
190
191pub trait ParachainBlockImportMarker {}
193
194impl<B: BlockT, BI, BE> ParachainBlockImportMarker for ParachainBlockImport<B, BI, BE> {}
195
196pub fn get_relay_slot(relay_header: &PHeader) -> Option<Slot> {
198 match sc_consensus_babe::find_pre_digest::<PBlock>(relay_header) {
199 Ok(pre_digest) => Some(pre_digest.slot()),
200 Err(err) => {
201 tracing::error!(
202 target: LOG_TARGET,
203 hash = %relay_header.hash(),
204 ?err,
205 "Relay chain block does not contain a BABE pre-digest. This should never happen.",
206 );
207 None
208 },
209 }
210}
211
212pub fn get_relay_slot_and_timestamp(
214 relay_header: &PHeader,
215 relay_slot_duration: Duration,
216) -> Option<(Slot, Timestamp)> {
217 get_relay_slot(relay_header).map(|slot| {
218 let t = Timestamp::new(relay_slot_duration.as_millis() as u64 * *slot);
219 (slot, t)
220 })
221}
222
223pub async fn load_abridged_host_configuration(
225 relay_parent: PHash,
226 relay_client: &impl RelayChainInterface,
227) -> Result<Option<AbridgedHostConfiguration>, RelayChainError> {
228 relay_client
229 .get_storage_by_key(relay_parent, relay_chain::well_known_keys::ACTIVE_CONFIG)
230 .await?
231 .map(|bytes| {
232 AbridgedHostConfiguration::decode(&mut &bytes[..])
233 .map_err(RelayChainError::DeserializationError)
234 })
235 .transpose()
236}