use crate::{
Client, Sub,
block::{self, Block},
};
use avail_rust_core::{H256, HasHeader};
use codec::Decode;
use std::{marker::PhantomData, time::Duration};
#[derive(Debug, Clone)]
pub struct ExtrinsicSubValue<T: HasHeader + Decode> {
pub list: Vec<block::BlockExtrinsic<T>>,
pub block_height: u32,
pub block_hash: H256,
}
#[derive(Clone)]
pub struct ExtrinsicSub<T: HasHeader + Decode> {
sub: Sub,
opts: block::extrinsic_options::Options,
_phantom: PhantomData<T>,
}
impl<T: HasHeader + Decode> ExtrinsicSub<T> {
pub fn new(client: Client, opts: block::extrinsic_options::Options) -> Self {
Self { sub: Sub::new(client), opts, _phantom: Default::default() }
}
pub async fn next(&mut self) -> Result<ExtrinsicSubValue<T>, crate::Error> {
loop {
let info = self.sub.next().await?;
let mut block = Block::new(self.sub.client_ref().clone(), info.hash).extrinsics();
block.set_retry_on_error(Some(self.sub.should_retry_on_error()));
let extrinsics = match block.all::<T>(self.opts.clone()).await {
Ok(x) => x,
Err(err) => {
self.sub.set_block_height(info.height);
return Err(err);
},
};
if extrinsics.is_empty() {
continue;
}
return Ok(ExtrinsicSubValue {
list: extrinsics,
block_hash: info.hash,
block_height: info.height,
});
}
}
pub fn set_opts(&mut self, value: block::extrinsic_options::Options) {
self.opts = value;
}
pub fn use_best_block(&mut self, value: bool) {
self.sub.use_best_block(value);
}
pub fn set_block_height(&mut self, block_height: u32) {
self.sub.set_block_height(block_height);
}
pub fn set_pool_rate(&mut self, value: Duration) {
self.sub.set_pool_rate(value);
}
pub fn set_retry_on_error(&mut self, value: Option<bool>) {
self.sub.set_retry_on_error(value);
}
pub fn should_retry_on_error(&self) -> bool {
self.sub.should_retry_on_error()
}
}
#[derive(Debug, Clone)]
pub struct EncodedExtrinsicSubValue {
pub list: Vec<block::BlockEncodedExtrinsic>,
pub block_height: u32,
pub block_hash: H256,
}
#[derive(Clone)]
pub struct EncodedExtrinsicSub {
sub: Sub,
opts: block::extrinsic_options::Options,
}
impl EncodedExtrinsicSub {
pub fn new(client: Client, opts: block::extrinsic_options::Options) -> Self {
Self { sub: Sub::new(client), opts }
}
pub async fn next(&mut self) -> Result<EncodedExtrinsicSubValue, crate::Error> {
loop {
let extrinsics = self.next_step().await?;
if extrinsics.list.is_empty() {
continue;
}
return Ok(extrinsics);
}
}
pub async fn next_step(&mut self) -> Result<EncodedExtrinsicSubValue, crate::Error> {
let info = self.sub.next().await?;
let mut block = Block::new(self.sub.client_ref().clone(), info.hash).encoded();
block.set_retry_on_error(Some(self.sub.should_retry_on_error()));
let extrinsics = match block.all(self.opts.clone()).await {
Ok(x) => x,
Err(err) => {
self.sub.set_block_height(info.height);
return Err(err);
},
};
return Ok(EncodedExtrinsicSubValue {
list: extrinsics,
block_hash: info.hash,
block_height: info.height,
});
}
pub fn set_opts(&mut self, value: block::extrinsic_options::Options) {
self.opts = value;
}
pub fn use_best_block(&mut self, value: bool) {
self.sub.use_best_block(value);
}
pub fn set_block_height(&mut self, block_height: u32) {
self.sub.set_block_height(block_height);
}
pub fn set_pool_rate(&mut self, value: Duration) {
self.sub.set_pool_rate(value);
}
pub fn set_retry_on_error(&mut self, value: Option<bool>) {
self.sub.set_retry_on_error(value);
}
pub fn should_retry_on_error(&self) -> bool {
self.sub.should_retry_on_error()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{clients::mock_client::MockClient, error::Error, prelude::*, subxt_rpcs::RpcClient};
use avail_rust_core::{
avail::data_availability::tx::SubmitData, rpc::system::fetch_extrinsics::ExtrinsicInformation,
};
#[tokio::test]
async fn extrinsic_sub_test() -> Result<(), Error> {
let (rpc_client, mut commander) = MockClient::new(TURING_ENDPOINT);
let client = Client::from_rpc_client(RpcClient::new(rpc_client)).await?;
let mut sub = ExtrinsicSub::<SubmitData>::new(client.clone(), Default::default());
sub.set_block_height(2326671);
let result = sub.next().await?;
assert_eq!(result.block_height, 2326672);
assert_eq!(result.list.len(), 1);
let result = sub.next().await?;
assert_eq!(result.block_height, 2326674);
assert_eq!(result.list.len(), 1);
sub.set_block_height(1);
assert_eq!(sub.sub.as_finalized().next_block_height, 1);
let mut data = ExtrinsicInformation::default();
let tx = client.tx().data_availability().submit_data("1234");
data.encoded = Some(const_hex::encode(tx.sign(&alice(), Options::new(2)).await?.encode()));
commander.extrinsics_ok(vec![data.clone()]); commander.extrinsics_ok(vec![]); commander.extrinsics_ok(vec![data.clone()]); commander.extrinsics_err(None); commander.extrinsics_ok(vec![data.clone()]);
let _ = sub.next().await?;
assert_eq!(sub.sub.as_finalized().next_block_height, 2);
let _ = sub.next().await?;
assert_eq!(sub.sub.as_finalized().next_block_height, 4);
sub.set_retry_on_error(Some(false));
let _ = sub.next().await.expect_err("Expect Error");
assert_eq!(sub.sub.as_finalized().next_block_height, 4);
let _ = sub.next().await?;
assert_eq!(sub.sub.as_finalized().next_block_height, 5);
Ok(())
}
#[tokio::test]
async fn raw_extrinsic_sub_test() -> Result<(), Error> {
let (rpc_client, mut commander) = MockClient::new(TURING_ENDPOINT);
let client = Client::from_rpc_client(RpcClient::new(rpc_client)).await?;
let opts = block::extrinsic_options::Options::new().filter((29u8, 1u8));
let mut sub = EncodedExtrinsicSub::new(client.clone(), opts);
sub.set_block_height(2326671);
let result = sub.next().await?;
assert_eq!(result.block_height, 2326672);
assert_eq!(result.list.len(), 1);
let result = sub.next().await?;
assert_eq!(result.block_height, 2326674);
assert_eq!(result.list.len(), 1);
sub.set_block_height(1);
assert_eq!(sub.sub.as_finalized().next_block_height, 1);
let mut data = ExtrinsicInformation::default();
let tx = client.tx().data_availability().submit_data("1234");
data.encoded = Some(const_hex::encode(tx.sign(&alice(), Options::new(2)).await?.encode()));
commander.extrinsics_ok(vec![data.clone()]); commander.extrinsics_ok(vec![]); commander.extrinsics_ok(vec![data.clone()]); commander.extrinsics_err(None); commander.extrinsics_ok(vec![data.clone()]);
let _ = sub.next().await?;
assert_eq!(sub.sub.as_finalized().next_block_height, 2);
let _ = sub.next().await?;
assert_eq!(sub.sub.as_finalized().next_block_height, 4);
sub.set_retry_on_error(Some(false));
let _ = sub.next().await.expect_err("Expect Error");
assert_eq!(sub.sub.as_finalized().next_block_height, 4);
let _ = sub.next().await?;
assert_eq!(sub.sub.as_finalized().next_block_height, 5);
Ok(())
}
}