From 6679c06303b13bdb538b225352e9411e4741a92e Mon Sep 17 00:00:00 2001 From: Jay White Date: Mon, 3 Aug 2026 03:12:02 -0400 Subject: [PATCH] feat: use client for `prover_db_indexer` --- pallets/prover_db_indexer/src/db_events.rs | 4 +- pallets/prover_db_indexer/src/lib.rs | 124 +++++----- .../src/offchain_consumer.rs | 49 ---- pallets/prover_db_indexer/src/tests.rs | 233 +++++++----------- 4 files changed, 145 insertions(+), 265 deletions(-) delete mode 100644 pallets/prover_db_indexer/src/offchain_consumer.rs diff --git a/pallets/prover_db_indexer/src/db_events.rs b/pallets/prover_db_indexer/src/db_events.rs index 21686d17..9e3ad054 100644 --- a/pallets/prover_db_indexer/src/db_events.rs +++ b/pallets/prover_db_indexer/src/db_events.rs @@ -37,7 +37,8 @@ pub enum DBEventError { } type RuntimeEvent = ::RuntimeEvent; -type EventRecord = frame_system::EventRecord, ::Hash>; +pub type EventRecord = + frame_system::EventRecord, ::Hash>; #[storage_alias] type Events = StorageValue, Vec>>; @@ -47,6 +48,7 @@ pub enum DBEvent { /// Table definitions have been updated. SchemaUpdated(Option, UpdateTableList), /// A table has been successfully dropped. + #[allow(dead_code, reason = "We mirror the pallet_tables TableDropped event")] TableDropped(Option, TableType, TableIdentifier, Source), /// This event is emitted when a quorum is reached amongst submissions and the /// data is finalized. diff --git a/pallets/prover_db_indexer/src/lib.rs b/pallets/prover_db_indexer/src/lib.rs index e0addc63..5fc8d9a6 100644 --- a/pallets/prover_db_indexer/src/lib.rs +++ b/pallets/prover_db_indexer/src/lib.rs @@ -49,10 +49,8 @@ mod mock; #[cfg(test)] mod tests; -#[expect(dead_code, reason = "Usage for this function is not yet implemented")] mod db_events; mod http_client; -mod offchain_consumer; /// Generated protobuf types for the indexer HTTP adapter wire format. mod proto { @@ -125,6 +123,12 @@ pub enum ConsumerError { /// inconsistency in the node's client state. #[snafu(display("finalized block hash mismatch"))] FinalizedBlockHashMismatch, + /// Querying the node's client for a block's captured events failed. + #[snafu(transparent)] + DBEvents { + /// Underlying error from [`crate::db_events::db_events_at`]. + source: crate::db_events::DBEventError, + }, } #[polkadot_sdk::frame_support::pallet] @@ -149,20 +153,28 @@ pub mod pallet { BlockEvent, EventCapture, ProverDbConsumerConfig, - TableIdentifierFilter, }; use crate::consumer_error::*; + use crate::db_events::{db_events_at, DBEvent, EventRecord}; use crate::ConsumerError; #[pallet::pallet] pub struct Pallet(_); #[pallet::config] - pub trait Config: polkadot_sdk::frame_system::Config {} + pub trait Config: + polkadot_sdk::frame_system::Config + + pallet_tables::Config + + pallet_indexing::Config + { + } #[pallet::hooks] - impl Hooks> for Pallet { + impl Hooks> for Pallet + where + EventRecord: TryInto>, + { // Fires at chain tip only (not during sync). Drains the offchain // DB queue, forwards to the HTTP server, deletes consumed entries. fn offchain_worker(_block_number: BlockNumberFor) { @@ -205,7 +217,10 @@ pub mod pallet { } } - impl Pallet { + impl Pallet + where + EventRecord: TryInto>, + { // ═══════════════════════════════════════════════════════════════ // CONSUMER: drain offchain DB → HTTP → delete (called from OCW) // ═══════════════════════════════════════════════════════════════ @@ -259,7 +274,7 @@ pub mod pallet { .skip(1) .take(config.max_blocks_per_invocation) { - Self::forward_block(&config.url, block_num, &config.include)?; + Self::forward_block(block_num, &config)?; // Checkpoint on the server (always, even for empty blocks). crate::http_client::checkpoint(&config.url, block_num).context(CheckpointSnafu)?; @@ -268,74 +283,45 @@ pub mod pallet { Ok(()) } - /// Forward a single block's events in extrinsic-index order, then - /// clear the offchain entries we consumed. If the block had no - /// captures (no high-water-mark key), this is a no-op. + /// Forward a single block's events. fn forward_block( - url: &url::Url, - block_num: u64, - include_filters: &[TableIdentifierFilter], - ) -> Result<(), ConsumerError> { - let Some(high_water_mark) = crate::offchain_consumer::read_high_water(block_num) else { - return Ok(()); - }; - - polkadot_sdk::sp_tracing::info!( - target: "prover_db_indexer", - "block {} — high-water-mark {}; probing for captured events", - block_num, - high_water_mark, - ); - - for ext_idx in 0..=high_water_mark { - let Some(events) = crate::offchain_consumer::read_events(block_num, ext_idx) else { - continue; - }; - Self::forward_events(url, block_num, &events, include_filters)?; - crate::offchain_consumer::clear_events(block_num, ext_idx); - } - - crate::offchain_consumer::clear_high_water(block_num); - Ok(()) - } - - /// POST one extrinsic's captured events to the indexer in deposit - /// order, skipping any whose table doesn't match this node's - /// include set. The capture queue is unfiltered (every validator - /// records the full block), so the filter applies only to the - /// HTTP forwarding done by this node's OCW. - fn forward_events( - url: &url::Url, block_num: u64, - events: &[BlockEvent<'_>], - include_filters: &[TableIdentifierFilter], + config: &ProverDbConsumerConfig, ) -> Result<(), ConsumerError> { - let filtered_events = events - .iter() - .filter(|event| table_matches_filters(event.table(), include_filters)); - for event in filtered_events { + let bn: BlockNumberFor = + block_num.checked_into().context(BlockNumberOverflowSnafu)?; + for event in db_events_at::(frame_system::Pallet::::block_hash(bn))? { match event { - BlockEvent::Drop(ident) => { - crate::http_client::drop_table(url, block_num, ident) - .context(DropTableSnafu)?; + DBEvent::TableDropped(_, _, table, _) => { + if table_matches_filters(&table, &config.include) { + crate::http_client::drop_table(&config.url, block_num, &table) + .context(DropTableSnafu)?; + } } - BlockEvent::Create(entry) => { - crate::http_client::create_table( - url, - block_num, - &entry.ident, - entry.ddl.to_vec(), - commitment_sql::ROW_NUMBER_COLUMN_NAME.into(), - ) - .context(CreateTableSnafu)?; + DBEvent::SchemaUpdated(_, updates) => { + updates + .into_iter() + .filter(|update| table_matches_filters(&update.ident, &config.include)) + .try_for_each(|update| { + crate::http_client::create_table( + &config.url, + block_num, + &update.ident, + update.create_statement.into_inner(), + commitment_sql::ROW_NUMBER_COLUMN_NAME.into(), + ) + .context(CreateTableSnafu) + })?; } - BlockEvent::Insert(entry) => { - crate::http_client::put_batches( - url, - block_num, - alloc::vec![(entry.table.as_ref(), entry.data.to_vec())], - ) - .context(PutBatchesSnafu)?; + DBEvent::QuorumReached { quorum, data } => { + if table_matches_filters(&quorum.table, &config.include) { + crate::http_client::put_batches( + &config.url, + block_num, + alloc::vec![(&quorum.table, data.into_inner())], + ) + .context(PutBatchesSnafu)?; + } } } } diff --git a/pallets/prover_db_indexer/src/offchain_consumer.rs b/pallets/prover_db_indexer/src/offchain_consumer.rs deleted file mode 100644 index 77eae0a6..00000000 --- a/pallets/prover_db_indexer/src/offchain_consumer.rs +++ /dev/null @@ -1,49 +0,0 @@ -//! OCW-only helpers for reading and clearing the offchain DB entries -//! the producer wrote. Producer-side helpers (key construction, types) -//! live in `sxt_core::prover_db_indexer`; this module wraps the OCW host -//! functions (`local_storage_get` / `local_storage_clear`) so the -//! consumer code in `lib.rs` doesn't repeat the boilerplate. - -use alloc::vec::Vec; - -use codec::Decode; -use polkadot_sdk::sp_core::offchain::StorageKind; -use sxt_core::prover_db_indexer::{key_for_event, key_for_high_water, BlockEvent}; - -/// Read the per-block high-water-mark. `None` means the block had no -/// captured events. -pub fn read_high_water(block: u64) -> Option { - let raw = polkadot_sdk::sp_io::offchain::local_storage_get( - StorageKind::PERSISTENT, - &key_for_high_water(block), - )?; - u32::decode(&mut &raw[..]).ok() -} - -/// Read the events emitted by a single extrinsic in a given block. -/// `None` means that extrinsic did not call `EventCapture::capture_events`. -/// The returned `BlockEvent<'static>` decodes into `Cow::Owned` for every -/// field, so the lifetime is purely a type-level annotation. -pub fn read_events(block: u64, extrinsic_index: u32) -> Option>> { - let raw = polkadot_sdk::sp_io::offchain::local_storage_get( - StorageKind::PERSISTENT, - &key_for_event(block, extrinsic_index), - )?; - Vec::>::decode(&mut &raw[..]).ok() -} - -/// Delete the per-block high-water-mark. -pub fn clear_high_water(block: u64) { - polkadot_sdk::sp_io::offchain::local_storage_clear( - StorageKind::PERSISTENT, - &key_for_high_water(block), - ); -} - -/// Delete the per-extrinsic event payload. -pub fn clear_events(block: u64, extrinsic_index: u32) { - polkadot_sdk::sp_io::offchain::local_storage_clear( - StorageKind::PERSISTENT, - &key_for_event(block, extrinsic_index), - ); -} diff --git a/pallets/prover_db_indexer/src/tests.rs b/pallets/prover_db_indexer/src/tests.rs index 948ef0aa..9d63d25d 100644 --- a/pallets/prover_db_indexer/src/tests.rs +++ b/pallets/prover_db_indexer/src/tests.rs @@ -10,16 +10,22 @@ use std::borrow::Cow; use codec::Encode; +use native_api::Api; +use pallet_tables::{CommitmentCreationCmd, UpdateTable}; use polkadot_sdk::frame_support::traits::Hooks; +use polkadot_sdk::frame_support::BoundedVec; +use polkadot_sdk::frame_system::EventRecord; use polkadot_sdk::sp_core::offchain::testing::{PendingRequest, TestOffchainExt}; -use polkadot_sdk::sp_core::offchain::{OffchainDbExt, OffchainStorage, OffchainWorkerExt}; +use polkadot_sdk::sp_core::offchain::{OffchainDbExt, OffchainWorkerExt}; +use polkadot_sdk::sp_core::storage::StorageData; use polkadot_sdk::sp_core::H256; use polkadot_sdk::sp_runtime::offchain::storage_lock::{StorageLock, Time}; use polkadot_sdk::sp_runtime::offchain::Duration; +use proof_of_sql_commitment_map::CommitmentSchemeFlags; use prost::Message; +use sxt_core::indexing::{BatchId, DataQuorum, SubmitterList}; use sxt_core::prover_db_indexer::{ key_for_event, - key_for_high_water, BlockEvent, CreateEntry, EventCapture, @@ -27,7 +33,7 @@ use sxt_core::prover_db_indexer::{ PROVER_DB_CONFIG_INCLUDE_KEY, PROVER_DB_CONFIG_URL_KEY, }; -use sxt_core::tables::TableIdentifier; +use sxt_core::tables::{QuorumScope, Source, TableIdentifier, TableType}; use crate::mock::*; use crate::proto; @@ -66,14 +72,18 @@ fn no_checkpoint_response() -> Vec { } /// Set a `*.*` filter. -fn setup_with_url(finalized_block_num: u32) -> (polkadot_sdk::sp_io::TestExternalities, StateArc) { - setup_with_config("*.*", finalized_block_num) +fn setup_with_url( + finalized_block_num: u32, + events: impl IntoIterator>>, +) -> (polkadot_sdk::sp_io::TestExternalities, StateArc) { + setup_with_config("*.*", finalized_block_num, events) } /// Set a non-empty include set. Used by the consumer-side filter tests. fn setup_with_config( filters: &str, finalized_block_num: u32, + events: impl IntoIterator>>, ) -> (polkadot_sdk::sp_io::TestExternalities, StateArc) { let mut ext = new_test_ext(); let (offchain, state) = TestOffchainExt::new(); @@ -81,7 +91,9 @@ fn setup_with_config( ext.register_extension(OffchainDbExt::new(offchain)); ext.register_extension(MockClientProvider::client_ext( Some((H256::zero(), finalized_block_num)), - [], + events + .into_iter() + .map(|e| Ok(Some(StorageData(e.encode())))), )); let mut config_store = std::collections::HashMap::new(); config_store.insert(PROVER_DB_CONFIG_URL_KEY.to_string(), MOCK_URL.to_string()); @@ -93,14 +105,50 @@ fn setup_with_config( (ext, state) } -/// Mirror what the producer would write for a block: one event payload -/// at `(block, ext_idx)` and the matching high-water-mark. -fn seed_block_events(state: &StateArc, block: u64, ext_idx: u32, events: Vec>) { - let mut s = state.write(); - s.persistent_storage - .set(b"", &key_for_event(block, ext_idx), &events.encode()); - s.persistent_storage - .set(b"", &key_for_high_water(block), &Encode::encode(&ext_idx)); +fn schema_updated_event( + name_namespace_ddl_tuples: impl IntoIterator)>, +) -> EventRecord { + event_record(pallet_tables::Event::::SchemaUpdated( + None, + BoundedVec::truncate_from( + name_namespace_ddl_tuples + .into_iter() + .map(|(name, namespace, ddl)| UpdateTable { + ident: TableIdentifier::from_str_unchecked(name, namespace), + create_statement: BoundedVec::truncate_from(ddl), + table_type: TableType::default(), + commitment: CommitmentCreationCmd::Empty(CommitmentSchemeFlags::default()), + source: Source::default(), + }) + .collect(), + ), + )) +} +fn table_dropped_event(name: &str, namespace: &str) -> EventRecord { + event_record(pallet_tables::Event::::TableDropped( + None, + TableType::default(), + TableIdentifier::from_str_unchecked(name, namespace), + Source::default(), + )) +} +fn quorum_reached_event( + name: &str, + namespace: &str, + data: Vec, +) -> EventRecord { + event_record(pallet_indexing::Event::::QuorumReached { + quorum: DataQuorum { + table: TableIdentifier::from_str_unchecked(name, namespace), + batch_id: BatchId::default(), + data_hash: H256::zero(), + block_number: Default::default(), + agreements: SubmitterList::default(), + dissents: SubmitterList::default(), + quorum_scope: QuorumScope::Public, + }, + data: BoundedVec::truncate_from(data), + }) } // ─── Tests ────────────────────────────────────────────────────────────── @@ -123,7 +171,7 @@ fn ocw_skips_when_not_configured() { /// the absence of queued expectations is the assertion. #[test] fn ocw_skips_when_lock_is_held() { - let (mut ext, _state) = setup_with_url(1); + let (mut ext, _state) = setup_with_url(1, [vec![]]); ext.execute_with(|| { let mut lock = StorageLock::