diff --git a/subxt/src/backend.rs b/subxt/src/backend.rs index e8504499764..bdef6dc1d41 100644 --- a/subxt/src/backend.rs +++ b/subxt/src/backend.rs @@ -85,6 +85,9 @@ pub trait Backend: 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>, 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>, BackendError>; + /// A stream of all new block headers as they arrive. async fn stream_all_block_headers( &self, diff --git a/subxt/src/backend/archive.rs b/subxt/src/backend/archive.rs index 712026b7886..080c6b1108b 100644 --- a/subxt/src/backend/archive.rs +++ b/subxt/src/backend/archive.rs @@ -180,6 +180,12 @@ impl Backend for ArchiveBackend { .await } + async fn latest_best_block_ref(&self) -> Result>, BackendError> { + Err(BackendError::Other( + "The archive backend cannot report the best block".into(), + )) + } + async fn stream_all_block_headers( &self, _hasher: T::Hasher, diff --git a/subxt/src/backend/chain_head.rs b/subxt/src/backend/chain_head.rs index 609b7d8b029..6313e41d83f 100644 --- a/subxt/src/backend/chain_head.rs +++ b/subxt/src/backend/chain_head.rs @@ -440,6 +440,15 @@ impl Backend for ChainHeadBackend { next_ref.ok_or_else(|| RpcError::SubscriptionDropped.into()) } + async fn latest_best_block_ref(&self) -> Result>, 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, diff --git a/subxt/src/backend/chain_head/follow_stream_driver.rs b/subxt/src/backend/chain_head/follow_stream_driver.rs index 5d0f150ccac..97539328aed 100644 --- a/subxt/src/backend/chain_head/follow_stream_driver.rs +++ b/subxt/src/backend/chain_head/follow_stream_driver.rs @@ -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; @@ -137,6 +138,30 @@ impl FollowStreamDriverSubscription { } } + /// 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> { + '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>> + Send + Sync { self.filter_map(|ev| std::future::ready(ev.into_event())) @@ -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( diff --git a/subxt/src/backend/combined.rs b/subxt/src/backend/combined.rs index 5877af6342c..6de8f8c3f2d 100644 --- a/subxt/src/backend/combined.rs +++ b/subxt/src/backend/combined.rs @@ -359,6 +359,18 @@ impl Backend for CombinedBackend { .await } + async fn latest_best_block_ref(&self) -> Result>, BackendError> { + try_backends( + &[ + // Ignore archive backend; it doesn't support this. + self.chainhead(), + self.legacy(), + ], + async |b: &dyn Backend| b.latest_best_block_ref().await, + ) + .await + } + async fn stream_all_block_headers( &self, hasher: T::Hasher, diff --git a/subxt/src/backend/legacy.rs b/subxt/src/backend/legacy.rs index 7c8ba91af89..7c93edcfeba 100644 --- a/subxt/src/backend/legacy.rs +++ b/subxt/src/backend/legacy.rs @@ -220,6 +220,18 @@ impl Backend for LegacyBackend { .await } + async fn latest_best_block_ref(&self) -> Result>, 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, diff --git a/subxt/src/client/online_client.rs b/subxt/src/client/online_client.rs index 6bf57625a68..8e311a3fed1 100644 --- a/subxt/src/client/online_client.rs +++ b/subxt/src/client/online_client.rs @@ -281,6 +281,21 @@ impl OnlineClient { 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>, 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, diff --git a/subxt/src/introduction.rs b/subxt/src/introduction.rs index 490a67c70a9..5aa8fd82485 100644 --- a/subxt/src/introduction.rs +++ b/subxt/src/introduction.rs @@ -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. //! diff --git a/testing/integration-tests/src/full_client/blocks.rs b/testing/integration-tests/src/full_client/blocks.rs index afb5e7d43b3..c0b93df03d0 100644 --- a/testing/integration-tests/src/full_client/blocks.rs +++ b/testing/integration-tests/src/full_client/blocks.rs @@ -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] diff --git a/testing/integration-tests/src/light_client.rs b/testing/integration-tests/src/light_client.rs index 638ea25ddb7..4f4f43eb5a9 100644 --- a/testing/integration-tests/src/light_client.rs +++ b/testing/integration-tests/src/light_client.rs @@ -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(); @@ -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;