From 61a126300aff21e66d27f8b31d3e82751696ce16 Mon Sep 17 00:00:00 2001 From: Vedhavyas Singareddi <7549475+vedhavyas@users.noreply.github.com> Date: Fri, 31 Jul 2026 08:30:15 +0530 Subject: [PATCH] decode block events with the metadata at that block Read the runtime version at each block and decode its events with the metadata that block ran, caching one metadata per spec version. Runtime upgrades insert event variants and shift the indices of later ones, so head metadata mis-reads historical events and drops everything after the first shifted variant in a block. --- indexer/src/staking.rs | 12 ++--- shared/src/subspace.rs | 120 ++++++++++++++++++++++++++++++++++++----- shared/src/types.rs | 9 ++++ 3 files changed, 123 insertions(+), 18 deletions(-) diff --git a/indexer/src/staking.rs b/indexer/src/staking.rs index f61a50c..f33bd14 100644 --- a/indexer/src/staking.rs +++ b/indexer/src/staking.rs @@ -452,10 +452,10 @@ mod tests { let block_ext = subspace.block_ext_at_hash(block_hash).await.unwrap(); let events = block_ext.events().await.unwrap(); - let completed: Vec<_> = events + let completed = events .find::() - .filter_map(|e| e.ok()) - .collect(); + .collect::, _>>() + .unwrap(); assert!( !completed.is_empty(), "should find DomainEpochCompleted event" @@ -563,10 +563,10 @@ mod tests { // Decode the event to learn the domain_id let events = block_ext.events().await.unwrap(); - let completed: Vec<_> = events + let completed = events .find::() - .filter_map(|e| e.ok()) - .collect(); + .collect::, _>>() + .unwrap(); assert!(!completed.is_empty()); let domain_id = completed[0].domain_id.clone(); diff --git a/shared/src/subspace.rs b/shared/src/subspace.rs index 295cf1d..d6109e4 100644 --- a/shared/src/subspace.rs +++ b/shared/src/subspace.rs @@ -1,6 +1,6 @@ //! Module that follows the chain and broadcast blocks data. use crate::error::Error; -use crate::types::{EventSegmentSize, system_event_segment_key}; +use crate::types::{EventSegmentSize, system_event_segment_key, system_events_key}; use futures_util::stream::Fuse; use futures_util::{StreamExt, TryStreamExt, stream}; use log::{debug, error, info, warn}; @@ -9,7 +9,7 @@ use sp_runtime::app_crypto::sp_core::crypto::Ss58AddressFormat; use sp_runtime::codec::{Decode, Encode, Input}; use sp_runtime::traits::{BlakeTwo256, Block as BlockT, Header as HeaderT}; use sp_runtime::{OpaqueExtrinsic, generic}; -use std::collections::{BTreeMap, BTreeSet}; +use std::collections::{BTreeMap, BTreeSet, HashMap}; use std::sync::Arc; use std::time::Duration; use subxt::backend::BackendExt; @@ -20,6 +20,7 @@ use subxt::{OnlineClient, SubstrateConfig}; use subxt_core::Config; use subxt_core::storage::address::StorageKey; use subxt_rpcs::{LegacyRpcMethods, RpcClient}; +use tokio::sync::RwLock; use tokio::sync::broadcast::{Receiver, Sender, channel}; /// Opaque block header type. @@ -50,6 +51,7 @@ pub type Nonce = u32; pub struct SubspaceBlockProvider { rpc: Arc, client: Arc, + metadata: MetadataCache, } impl SubspaceBlockProvider { @@ -92,6 +94,7 @@ impl SubspaceBlockProvider { state_root, extrinsics_root, client: self.client.clone(), + metadata: self.metadata.clone(), }) } } @@ -136,6 +139,7 @@ pub struct BlockExt { pub state_root: BlockHash, pub extrinsics_root: BlockHash, client: Arc, + metadata: MetadataCache, } impl BlockExt { @@ -252,8 +256,14 @@ impl BlockExt { /// Returns block events pub async fn events(&self) -> Result, Error> { - let events = self.client.events().at(self.hash).await?; - Ok(events) + let metadata = self.metadata.at(self.hash).await?; + let event_data = self + .client + .backend() + .storage_fetch_value(system_events_key(), self.hash) + .await? + .unwrap_or_default(); + Ok(Events::decode_from(event_data, metadata)) } pub async fn events_from_segments(&self) -> Result>, Error> { @@ -261,6 +271,7 @@ impl BlockExt { let event_count = self .read_storage::<_, u32>("System", "EventCount", ()) .await?; + let metadata = self.metadata.at(self.hash).await?; let max_segment = event_count / segment_size; let mut total_events = vec![]; let backend = self.client.backend(); @@ -268,10 +279,9 @@ impl BlockExt { let key = system_event_segment_key(segment); let data = backend.storage_fetch_value(key.clone(), self.hash).await?; if let Some(event_data) = data { - let events = - Events::::decode_from(event_data, self.client.metadata()) - .iter() - .collect::, _>>()?; + let events = Events::::decode_from(event_data, metadata.clone()) + .iter() + .collect::, _>>()?; total_events.extend(events); } } @@ -314,6 +324,7 @@ const STARTUP_CONNECT_TIMEOUT: Duration = Duration::from_secs(60); pub struct Subspace { rpc: Arc, client: Arc, + metadata: MetadataCache, sink: BlocksSink, stream: BlocksStream, } @@ -332,9 +343,12 @@ pub struct NetworkDetails { /// (V16/V15); on this runtime that metadata mis-sizes some events, so subxt /// overruns the SCALE buffer while iterating a block's events. V14 decodes /// correctly, so it is pinned on the client. -async fn fetch_stable_metadata(rpc: &SubspaceRpcClient) -> Result { +async fn fetch_stable_metadata( + rpc: &SubspaceRpcClient, + at: Option, +) -> Result { let prefixed = rpc - .state_get_metadata(None) + .state_get_metadata(at) .await? .to_frame_metadata() .map_err(|e| Error::Config(format!("decode stable metadata: {e}")))?; @@ -342,6 +356,42 @@ async fn fetch_stable_metadata(rpc: &SubspaceRpcClient) -> Result, + by_spec_version: Arc>>, +} + +impl MetadataCache { + fn new(rpc: Arc) -> Self { + Self { + rpc, + by_spec_version: Arc::default(), + } + } + + async fn at(&self, hash: BlockHash) -> Result { + let spec_version = self + .rpc + .state_get_runtime_version(Some(hash)) + .await? + .spec_version; + + if let Some(metadata) = self.by_spec_version.read().await.get(&spec_version) { + return Ok(metadata.clone()); + } + + let metadata = fetch_stable_metadata(&self.rpc, Some(hash)).await?; + self.by_spec_version + .write() + .await + .insert(spec_version, metadata.clone()); + Ok(metadata) + } +} + impl Subspace { pub async fn new_from_url(url: &str) -> Result { tokio::time::timeout(STARTUP_CONNECT_TIMEOUT, Self::connect(url)) @@ -359,9 +409,10 @@ impl Subspace { let rpc = Arc::new(LegacyRpcMethods::::new(rpc_client.clone())); let client = Arc::new(SubspaceClient::from_rpc_client(rpc_client).await?); // Override subxt's default (unstable) metadata with stable V14. - client.set_metadata(fetch_stable_metadata(&rpc).await?); + client.set_metadata(fetch_stable_metadata(&rpc, None).await?); let (sink, stream) = channel(100); Ok(Self { + metadata: MetadataCache::new(rpc.clone()), rpc, client, sink, @@ -383,7 +434,7 @@ impl Subspace { let mut updates = client.updater().runtime_updates().await?; while let Some(update) = updates.next().await { update?; - client.set_metadata(fetch_stable_metadata(&rpc).await?); + client.set_metadata(fetch_stable_metadata(&rpc, None).await?); } Ok(()) } @@ -397,6 +448,7 @@ impl Subspace { SubspaceBlockProvider { rpc: self.rpc.clone(), client: self.client.clone(), + metadata: self.metadata.clone(), } } @@ -595,6 +647,7 @@ impl Subspace { state_root, extrinsics_root, client: self.client.clone(), + metadata: self.metadata.clone(), }) } @@ -756,3 +809,46 @@ impl sp_blockchain::HeaderMetadata for HeadersMetadataCache { // nothing to do here } } + +#[cfg(test)] +mod tests { + use super::*; + use std::str::FromStr; + + const RPC_URL: &str = "wss://rpc.mainnet.autonomys.xyz/ws"; + + /// Block 6721910 ran spec 9, where `Balances::Burned` is variant 11; head metadata + /// (spec 11) has it at 12 and derails every event after it in this block. + #[tokio::test] + async fn test_events_decode_with_metadata_at_block() { + let subspace = Subspace::new_from_url(RPC_URL) + .await + .unwrap() + .block_provider(); + + let hash = BlockHash::from_str( + "0xb26dc651dd8317b593775da3202061dd0c1dea817e0e60c5f0f4b14c6f9efb39", + ) + .unwrap(); + let events = subspace + .block_ext_at_hash(hash) + .await + .unwrap() + .events() + .await + .unwrap(); + + let decoded = events.iter().collect::, _>>().unwrap(); + assert_eq!( + (decoded[1].pallet_name(), decoded[1].variant_name()), + ("Balances", "Burned") + ); + assert!( + decoded + .iter() + .any(|e| e.pallet_name() == "Domains" + && e.variant_name() == "DomainEpochCompleted"), + "should find DomainEpochCompleted event" + ); + } +} diff --git a/shared/src/types.rs b/shared/src/types.rs index afd9aa8..56df9ec 100644 --- a/shared/src/types.rs +++ b/shared/src/types.rs @@ -15,6 +15,15 @@ impl Address for EventSegmentSize { } } +pub(crate) fn system_events_key() -> Vec { + let a = sp_crypto_hashing::twox_128(b"System"); + let b = sp_crypto_hashing::twox_128(b"Events"); + let mut key = vec![]; + key.extend_from_slice(&a); + key.extend_from_slice(&b); + key +} + pub(crate) fn system_event_segment_key(segment: u32) -> Vec { let a = sp_crypto_hashing::twox_128(b"System"); let b = sp_crypto_hashing::twox_128(b"EventSegments");