Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions subxt/src/backend.rs
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,9 @@ pub trait Backend<T: Config>: sealed::Sealed + Send + Sync + 'static {
/// Note: needed only in blocks client for finalized block stream; can prolly be removed.
async fn latest_finalized_block_ref(&self) -> Result<BlockRef<HashFor<T>>, BackendError>;

/// Get the current best block hash. It is not finalized and may be reorged or pruned.
async fn latest_best_block_ref(&self) -> Result<BlockRef<HashFor<T>>, BackendError>;

/// A stream of all new block headers as they arrive.
async fn stream_all_block_headers(
&self,
Expand Down
6 changes: 6 additions & 0 deletions subxt/src/backend/archive.rs
Original file line number Diff line number Diff line change
Expand Up @@ -180,6 +180,12 @@ impl<T: Config> Backend<T> for ArchiveBackend<T> {
.await
}

async fn latest_best_block_ref(&self) -> Result<BlockRef<HashFor<T>>, BackendError> {
Err(BackendError::Other(
"The archive backend cannot report the best block".into(),
))
}

async fn stream_all_block_headers(
&self,
_hasher: T::Hasher,
Expand Down
9 changes: 9 additions & 0 deletions subxt/src/backend/chain_head.rs
Original file line number Diff line number Diff line change
Expand Up @@ -440,6 +440,15 @@ impl<T: Config> Backend<T> for ChainHeadBackend<T> {
next_ref.ok_or_else(|| RpcError::SubscriptionDropped.into())
}

async fn latest_best_block_ref(&self) -> Result<BlockRef<HashFor<T>>, BackendError> {
self.follow_handle
.subscribe()
.latest_best_block()
.await
.map(Into::into)
.ok_or_else(|| RpcError::SubscriptionDropped.into())
}

async fn stream_all_block_headers(
&self,
_hasher: T::Hasher,
Expand Down
108 changes: 108 additions & 0 deletions subxt/src/backend/chain_head/follow_stream_driver.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
use super::follow_stream_unpin::{BlockRef, FollowStreamMsg, FollowStreamUnpin};
use crate::config::Hash;
use crate::error::{BackendError, RpcError};
use futures::FutureExt;
use futures::stream::{Stream, StreamExt};
use std::collections::{HashMap, HashSet, VecDeque};
use std::ops::DerefMut;
Expand Down Expand Up @@ -137,6 +138,30 @@ impl<H: Hash> FollowStreamDriverSubscription<H> {
}
}

/// The last queued `BestBlockChanged` after `Initialized`, else the latest finalized block.
/// A queued `Stop` restarts the wait. `None` if the stream ends.
pub async fn latest_best_block(mut self) -> Option<BlockRef<H>> {
'restart: loop {
let mut best_block = loop {
if let FollowStreamMsg::Event(FollowEvent::Initialized(init)) = self.next().await? {
break init.finalized_block_hashes.last().cloned();
}
};

while let Some(Some(msg)) = self.next().now_or_never() {
match msg {
FollowStreamMsg::Event(FollowEvent::BestBlockChanged(ev)) => {
best_block = Some(ev.best_block_hash);
}
FollowStreamMsg::Event(FollowEvent::Stop) => continue 'restart,
_ => {}
}
}

return best_block;
}
}

/// Subscribe to the follow events, ignoring any other messages.
pub fn events(self) -> impl Stream<Item = FollowEvent<BlockRef<H>>> + Send + Sync {
self.filter_map(|ev| std::future::ready(ev.into_event()))
Expand Down Expand Up @@ -722,6 +747,89 @@ mod test {
assert_eq!(evs, expected);
}

#[tokio::test]
async fn latest_best_block_prefers_replayed_best_block_over_finalized() {
let mut driver = test_follow_stream_driver_getter(
|| {
[
Ok(ev_initialized(0)),
Ok(ev_new_block(0, 1)),
Ok(ev_best_block(1)),
Ok(ev_new_block(1, 2)),
Ok(ev_best_block(2)),
Err(BackendError::other("ended")),
]
},
10,
);

let _ready = driver.next().await.unwrap();
let _init0 = driver.next().await.unwrap();
let _new1 = driver.next().await.unwrap();
let _best1 = driver.next().await.unwrap();
let _new2 = driver.next().await.unwrap();
let _best2 = driver.next().await.unwrap();

let best_block = driver.handle().subscribe().latest_best_block().await;
assert_eq!(best_block.map(|b| b.hash()), Some(H256::from_low_u64_le(2)));
}

#[tokio::test]
async fn latest_best_block_falls_back_to_finalized_when_nothing_replayed() {
let mut driver = test_follow_stream_driver_getter(
|| {
[
Ok(ev_initialized(0)),
Ok(ev_new_block(0, 1)),
Ok(ev_best_block(1)),
Ok(ev_finalized([1], [])),
Err(BackendError::other("ended")),
]
},
10,
);

let _ready = driver.next().await.unwrap();
let _init0 = driver.next().await.unwrap();
let _new1 = driver.next().await.unwrap();
let _best1 = driver.next().await.unwrap();
let _fin1 = driver.next().await.unwrap();

let best_block = driver.handle().subscribe().latest_best_block().await;
assert_eq!(best_block.map(|b| b.hash()), Some(H256::from_low_u64_le(1)));
}

#[tokio::test]
async fn latest_best_block_restarts_after_a_stop_in_the_replay() {
let mut driver = test_follow_stream_driver_getter(
|| {
[
Ok(ev_initialized(0)),
Ok(ev_new_block(0, 1)),
Ok(ev_best_block(1)),
Ok(FollowEvent::Stop),
Ok(ev_initialized(2)),
Err(BackendError::other("ended")),
]
},
10,
);

let _ready = driver.next().await.unwrap();
let _init0 = driver.next().await.unwrap();
let _new1 = driver.next().await.unwrap();
let _best1 = driver.next().await.unwrap();

let subscription = driver.handle().subscribe();

let _stop = driver.next().await.unwrap();
let _ready_again = driver.next().await.unwrap();
let _init2 = driver.next().await.unwrap();

let best_block = subscription.latest_best_block().await;
assert_eq!(best_block.map(|b| b.hash()), Some(H256::from_low_u64_le(2)));
}

#[tokio::test]
async fn subscribe_finalized_blocks_restart_works() {
let mut driver = test_follow_stream_driver_getter(
Expand Down
12 changes: 12 additions & 0 deletions subxt/src/backend/combined.rs
Original file line number Diff line number Diff line change
Expand Up @@ -359,6 +359,18 @@ impl<T: Config> Backend<T> for CombinedBackend<T> {
.await
}

async fn latest_best_block_ref(&self) -> Result<BlockRef<HashFor<T>>, BackendError> {
try_backends(
&[
// Ignore archive backend; it doesn't support this.
self.chainhead(),
self.legacy(),
],
async |b: &dyn Backend<T>| b.latest_best_block_ref().await,
)
.await
}

async fn stream_all_block_headers(
&self,
hasher: T::Hasher,
Expand Down
12 changes: 12 additions & 0 deletions subxt/src/backend/legacy.rs
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,18 @@ impl<T: Config> Backend<T> for LegacyBackend<T> {
.await
}

async fn latest_best_block_ref(&self) -> Result<BlockRef<HashFor<T>>, BackendError> {
retry(|| async {
let hash = self
.methods
.chain_get_block_hash(None)
.await?
.ok_or_else(|| BackendError::other("chain_getBlockHash returned no best block"))?;
Ok(BlockRef::from_hash(hash))
})
.await
}

async fn stream_all_block_headers(
&self,
hasher: T::Hasher,
Expand Down
15 changes: 15 additions & 0 deletions subxt/src/client/online_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -281,6 +281,21 @@ impl<T: Config> OnlineClient<T> {
self.at_block(latest_block).await
}

/// Instantiate a client to work at the current best block _at the time of instantiation_.
/// This does not track new blocks. The best block is not finalized and may be reorged or pruned.
pub async fn at_current_best_block(
&self,
) -> Result<ClientAtBlock<T, OnlineClientAtBlockImpl<T>>, OnlineClientAtBlockError> {
let best_block = self
.inner
.backend
.latest_best_block_ref()
.await
.map_err(|e| OnlineClientAtBlockError::CannotGetCurrentBlock { reason: e })?;

self.at_block(best_block).await
}

/// Instantiate a client for working at a specific block.
pub async fn at_block(
&self,
Expand Down
6 changes: 3 additions & 3 deletions subxt/src/introduction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -79,9 +79,9 @@
//! Read the [`crate::config`] docs for more.
//! 2. Create a _client_ for interacting with the chain, which consumes this configuration.
//! Read the [`crate::client`] docs for more.
//! 3. Pick a block to work at. To work at the current block at the time of calling, you'd use
//! [`crate::client::OnlineClient::at_current_block()`]. To stream blocks, you can use
//! [`crate::client::OnlineClient::stream_blocks()`] and similar.
//! 3. Pick a block to work at. To work at the current finalized block at the time of calling, you'd use
//! [`crate::client::OnlineClient::at_current_block()`], or [`crate::client::OnlineClient::at_current_best_block()`]
//! for the best block. To stream blocks, you can use [`crate::client::OnlineClient::stream_blocks()`] and similar.
//! 4. Do things in the context of this block. See the examples for more, or explore the documentation starting at
//! [`crate::client::ClientAtBlock`] to dig into the various things you can do at a given block.
//!
Expand Down
25 changes: 25 additions & 0 deletions testing/integration-tests/src/full_client/blocks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,31 @@ async fn finalized_headers_subscription() -> Result<(), subxt::Error> {
Ok(())
}

#[subxt_test]
async fn best_block_is_at_or_ahead_of_finalized_block() -> Result<(), subxt::Error> {
let ctx = test_context().await;
let api = ctx.client();

let mut sub = api.stream_blocks().await?;
consume_initial_blocks(&mut sub).await;

for _ in 0..2 {
let finalized = sub.next().await.unwrap()?;
let best = api.at_current_best_block().await?;
assert!(
best.block_number() >= finalized.number(),
"best block {} is behind finalized block {}",
best.block_number(),
finalized.number()
);
if best.block_number() == finalized.number() {
assert_eq!(best.block_hash(), finalized.hash());
}
}

Ok(())
}

// This test only uses legacy RPCs; only run once for default backend + rpc client.
#[cfg(all(default_backend, default_rpc))]
#[subxt_test]
Expand Down
18 changes: 18 additions & 0 deletions testing/integration-tests/src/light_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,6 +105,23 @@ async fn non_finalized_headers_subscription(api: &Client) -> Result<(), subxt::E
Ok(())
}

// Check that the best block can be fetched and is not behind the finalized block.
async fn best_block_lookup(api: &Client) -> Result<(), subxt::Error> {
tracing::trace!("Check best_block_lookup");

let finalized = api.at_current_block().await?;
let best = api.at_current_best_block().await?;

assert!(
best.block_number() >= finalized.block_number(),
"best block {} is behind finalized block {}",
best.block_number(),
finalized.block_number()
);

Ok(())
}

// Check that we can subscribe to finalized blocks.
async fn finalized_headers_subscription(api: &Client) -> Result<(), subxt::Error> {
let now = std::time::Instant::now();
Expand Down Expand Up @@ -256,6 +273,7 @@ async fn light_client_tests() {
finalized_headers_subscription(&api),
)
.await;
run_check("best_block_lookup", best_block_lookup(&api)).await;
run_check("runtime_api_call", runtime_api_call(&api)).await;
run_check("storage_plain_lookup", storage_plain_lookup(&api)).await;
run_check("dynamic_constant_query", dynamic_constant_query(&api)).await;
Expand Down
Loading