mirror of
https://github.com/pezkuwichain/pezkuwi-subxt.git
synced 2026-07-22 13:45:40 +00:00
Batch transactions in complex relays (#1669)
* batch transactions in message relay: API prototype * get rid of Box<dyn BatchTransaction> and actually submit it * test batch transactions * message_lane_loop_works_with_batch_transactions * removed logger * BatchConfirmationTransaction + BatchDeliveryTransaction * more prototyping * fmt * continue with batch calls * impl BatchCallBuilder for () * BatchDeliveryTransaction impl * BundledBatchCallBuilder * proper impl of BundledBatchCallBuilder + use it in RialtoParachain -> Millau * impl prove_header in OnDemandHeadersRelay * impl OnDemandParachainsRelay::prove_header (needs extensive tests) * added a couple of TODOs * return Result<Option<BatchTx>> when asking for more headers * prove headers when reauire_* is called && return proper headers from required_header_id * split parachains::prove_header and test select_headers_to_prove * more traces and leave TODOs * use finality stream in SubstrateFinalitySource::prove_block_finality * prove parachain head at block, selected by headers relay * const ANCIENT_BLOCK_THRESHOLD * TODO -> proof * clippy and spelling * BatchCallBuilder::build_batch_call() returns Result * read first proof from two streams * FailedToFindFinalityProof -> FinalityProofNotFound * changed select_headers_to_prove to version from PR review
This commit is contained in:
committed by
Bastian Köcher
parent
a732a04ed4
commit
be27bd5e97
@@ -20,13 +20,18 @@ use crate::finality::{engine::Engine, FinalitySyncPipelineAdapter, SubstrateFina
|
||||
|
||||
use async_std::sync::{Arc, Mutex};
|
||||
use async_trait::async_trait;
|
||||
use bp_header_chain::FinalityProof;
|
||||
use codec::Decode;
|
||||
use finality_relay::SourceClient;
|
||||
use futures::stream::{unfold, Stream, StreamExt};
|
||||
use futures::{
|
||||
select,
|
||||
stream::{try_unfold, unfold, Stream, StreamExt, TryStreamExt},
|
||||
};
|
||||
use num_traits::One;
|
||||
use relay_substrate_client::{
|
||||
BlockNumberOf, BlockWithJustification, Chain, Client, Error, HeaderOf,
|
||||
};
|
||||
use relay_utils::relay_loop::Client as RelayClient;
|
||||
use relay_utils::{relay_loop::Client as RelayClient, UniqueSaturatedInto};
|
||||
use std::pin::Pin;
|
||||
|
||||
/// Shared updatable reference to the maximal header number that we want to sync from the source.
|
||||
@@ -70,6 +75,111 @@ impl<P: SubstrateFinalitySyncPipeline> SubstrateFinalitySource<P> {
|
||||
// target node may be missing proofs that are already available at the source
|
||||
self.client.best_finalized_header_number().await
|
||||
}
|
||||
|
||||
/// Return header and its justification of the given block or its descendant that
|
||||
/// has a GRANDPA justification.
|
||||
///
|
||||
/// This method is optimized for cases when `block_number` is close to the best finalized
|
||||
/// chain block.
|
||||
pub async fn prove_block_finality(
|
||||
&self,
|
||||
block_number: BlockNumberOf<P::SourceChain>,
|
||||
) -> Result<
|
||||
(relay_substrate_client::SyncHeader<HeaderOf<P::SourceChain>>, SubstrateFinalityProof<P>),
|
||||
Error,
|
||||
> {
|
||||
// first, subscribe to proofs
|
||||
let next_persistent_proof =
|
||||
self.persistent_proofs_stream(block_number + One::one()).await?.fuse();
|
||||
let next_ephemeral_proof = self.ephemeral_proofs_stream(block_number).await?.fuse();
|
||||
|
||||
// in perfect world we'll need to return justfication for the requested `block_number`
|
||||
let (header, maybe_proof) = self.header_and_finality_proof(block_number).await?;
|
||||
if let Some(proof) = maybe_proof {
|
||||
return Ok((header, proof))
|
||||
}
|
||||
|
||||
// otherwise we don't care which header to return, so let's select first
|
||||
futures::pin_mut!(next_persistent_proof, next_ephemeral_proof);
|
||||
loop {
|
||||
select! {
|
||||
maybe_header_and_proof = next_persistent_proof.next() => match maybe_header_and_proof {
|
||||
Some(header_and_proof) => return header_and_proof,
|
||||
None => continue,
|
||||
},
|
||||
maybe_header_and_proof = next_ephemeral_proof.next() => match maybe_header_and_proof {
|
||||
Some(header_and_proof) => return header_and_proof,
|
||||
None => continue,
|
||||
},
|
||||
complete => return Err(Error::FinalityProofNotFound(block_number.unique_saturated_into()))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Returns stream of headers and their persistent proofs, starting from given block.
|
||||
async fn persistent_proofs_stream(
|
||||
&self,
|
||||
block_number: BlockNumberOf<P::SourceChain>,
|
||||
) -> Result<
|
||||
impl Stream<
|
||||
Item = Result<
|
||||
(
|
||||
relay_substrate_client::SyncHeader<HeaderOf<P::SourceChain>>,
|
||||
SubstrateFinalityProof<P>,
|
||||
),
|
||||
Error,
|
||||
>,
|
||||
>,
|
||||
Error,
|
||||
> {
|
||||
let client = self.client.clone();
|
||||
let best_finalized_block_number = self.client.best_finalized_header_number().await?;
|
||||
Ok(try_unfold((client, block_number), move |(client, current_block_number)| async move {
|
||||
// if we've passed the `best_finalized_block_number`, we no longer need persistent
|
||||
// justifications
|
||||
if current_block_number > best_finalized_block_number {
|
||||
return Ok(None)
|
||||
}
|
||||
|
||||
let (header, maybe_proof) =
|
||||
header_and_finality_proof::<P>(&client, current_block_number).await?;
|
||||
let next_block_number = current_block_number + One::one();
|
||||
let next_state = (client, next_block_number);
|
||||
|
||||
Ok(Some((maybe_proof.map(|proof| (header, proof)), next_state)))
|
||||
})
|
||||
.try_filter_map(|maybe_result| async { Ok(maybe_result) }))
|
||||
}
|
||||
|
||||
/// Returns stream of headers and their ephemeral proofs, starting from given block.
|
||||
async fn ephemeral_proofs_stream(
|
||||
&self,
|
||||
block_number: BlockNumberOf<P::SourceChain>,
|
||||
) -> Result<
|
||||
impl Stream<
|
||||
Item = Result<
|
||||
(
|
||||
relay_substrate_client::SyncHeader<HeaderOf<P::SourceChain>>,
|
||||
SubstrateFinalityProof<P>,
|
||||
),
|
||||
Error,
|
||||
>,
|
||||
>,
|
||||
Error,
|
||||
> {
|
||||
let client = self.client.clone();
|
||||
Ok(self.finality_proofs().await?.map(Ok).try_filter_map(move |proof| {
|
||||
let client = client.clone();
|
||||
async move {
|
||||
if proof.target_header_number() < block_number {
|
||||
return Ok(None)
|
||||
}
|
||||
|
||||
let header = client.header_by_number(proof.target_header_number()).await?;
|
||||
Ok(Some((header.into(), proof)))
|
||||
}
|
||||
}))
|
||||
}
|
||||
}
|
||||
|
||||
impl<P: SubstrateFinalitySyncPipeline> Clone for SubstrateFinalitySource<P> {
|
||||
@@ -119,18 +229,7 @@ impl<P: SubstrateFinalitySyncPipeline> SourceClient<FinalitySyncPipelineAdapter<
|
||||
),
|
||||
Error,
|
||||
> {
|
||||
let header_hash = self.client.block_hash_by_number(number).await?;
|
||||
let signed_block = self.client.get_block(Some(header_hash)).await?;
|
||||
|
||||
let justification = signed_block
|
||||
.justification(P::FinalityEngine::ID)
|
||||
.map(|raw_justification| {
|
||||
SubstrateFinalityProof::<P>::decode(&mut raw_justification.as_slice())
|
||||
})
|
||||
.transpose()
|
||||
.map_err(Error::ResponseParseFailed)?;
|
||||
|
||||
Ok((signed_block.header().into(), justification))
|
||||
header_and_finality_proof::<P>(&self.client, number).await
|
||||
}
|
||||
|
||||
async fn finality_proofs(&self) -> Result<Self::FinalityProofsStream, Error> {
|
||||
@@ -173,3 +272,27 @@ impl<P: SubstrateFinalitySyncPipeline> SourceClient<FinalitySyncPipelineAdapter<
|
||||
.boxed())
|
||||
}
|
||||
}
|
||||
|
||||
async fn header_and_finality_proof<P: SubstrateFinalitySyncPipeline>(
|
||||
client: &Client<P::SourceChain>,
|
||||
number: BlockNumberOf<P::SourceChain>,
|
||||
) -> Result<
|
||||
(
|
||||
relay_substrate_client::SyncHeader<HeaderOf<P::SourceChain>>,
|
||||
Option<SubstrateFinalityProof<P>>,
|
||||
),
|
||||
Error,
|
||||
> {
|
||||
let header_hash = client.block_hash_by_number(number).await?;
|
||||
let signed_block = client.get_block(Some(header_hash)).await?;
|
||||
|
||||
let justification = signed_block
|
||||
.justification(P::FinalityEngine::ID)
|
||||
.map(|raw_justification| {
|
||||
SubstrateFinalityProof::<P>::decode(&mut raw_justification.as_slice())
|
||||
})
|
||||
.transpose()
|
||||
.map_err(Error::ResponseParseFailed)?;
|
||||
|
||||
Ok((signed_block.header().into(), justification))
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user