diff --git a/lez/indexer/core/src/block_store.rs b/lez/indexer/core/src/block_store.rs index 99eab96f0..62faa70a5 100644 --- a/lez/indexer/core/src/block_store.rs +++ b/lez/indexer/core/src/block_store.rs @@ -114,17 +114,8 @@ impl IndexerStore { Ok(self.dbio.get_block_events_range(from, to)?) } - pub fn get_events_by_tx_hash(&self, tx_hash: [u8; 32]) -> Result> { - let Some(block_id) = self.dbio.get_block_id_by_tx_hash(tx_hash)? else { - return Ok(None); - }; - let Some(groups) = self.dbio.get_block_events(block_id)? else { - return Ok(None); - }; - Ok(groups - .into_iter() - .find(|group| group.tx_hash.0 == tx_hash) - .map(|group| (block_id, group))) + pub fn block_id_by_tx_hash(&self, tx_hash: [u8; 32]) -> Result> { + Ok(self.dbio.get_block_id_by_tx_hash(tx_hash)?) } pub fn get_block_by_hash(&self, hash: [u8; 32]) -> Result> { @@ -452,7 +443,7 @@ mod tests { use testnet_initial_state::initial_pub_accounts_private_keys; use super::*; - use crate::event_filter::SelectorFilter; + use crate::event_filter::{SelectorFilter, covered_over_range}; // Host-side mirror of the `event_emitter` test guest's instruction. #[derive(serde::Serialize)] @@ -651,11 +642,16 @@ mod tests { vec![(3, groups.clone())] ); - let (block_id, group) = store - .get_events_by_tx_hash(invoke_hash.0) + let block_id = store + .block_id_by_tx_hash(invoke_hash.0) .unwrap() - .expect("tx-hash lookup must find the group"); + .expect("the invoking tx must resolve its block"); assert_eq!(block_id, 3); + let group = store + .get_events_for_block(block_id) + .unwrap() + .and_then(|rows| rows.into_iter().find(|row| row.tx_hash == invoke_hash)) + .expect("tx-hash lookup must find the group"); assert_eq!(group, groups[0]); } @@ -706,7 +702,13 @@ mod tests { // archival run stores both emitted events. assert_eq!(store.get_events_for_block(3).unwrap(), None); assert!(store.get_events_range(1, 3).unwrap().is_empty()); - assert_eq!(store.get_events_by_tx_hash(invoke_hash.0).unwrap(), None); + assert!( + store + .block_id_by_tx_hash(invoke_hash.0) + .unwrap() + .and_then(|block_id| store.get_events_for_block(block_id).unwrap()) + .is_none() + ); } #[tokio::test] @@ -806,6 +808,35 @@ mod tests { ); } + #[tokio::test] + async fn coverage_follows_the_segment_history_end_to_end() { + let home = tempdir().unwrap(); + let tip = { + let store = open_with(home.as_ref(), EventFilter::Archival); + seed_emitted_events(&store).await; + store.get_last_block_id().unwrap().unwrap() + }; + + let reopened = open_default(home.as_ref()); + let segments = reopened.filter_segments(); + + assert!(covered_over_range(segments, 1, tip, None, None)); + assert!(!covered_over_range( + segments, + 1, + tip.saturating_add(1), + None, + None + )); + assert!(!covered_over_range( + segments, + tip.saturating_add(1), + tip.saturating_add(1), + None, + None + )); + } + #[tokio::test] async fn sources_segment_survives_reopen_regardless_of_insertion_order() { let filter = |entries: Vec<(ProgramId, SelectorFilter)>| { @@ -933,6 +964,21 @@ mod tests { assert!(IndexerStore::open_db(home.as_ref(), Vec::new(), EventFilter::Archival).is_err()); } + + #[tokio::test] + async fn filtered_out_tx_still_resolves_its_block() { + let home = tempdir().unwrap(); + let store = open_default(home.as_ref()); + let invoke_hash = seed_emitted_events(&store).await; + + // The events row was dropped by the filter, but the height is still known — + // which is what lets the query layer reject instead of serving `[]`. + let block_id = store + .block_id_by_tx_hash(invoke_hash.0) + .unwrap() + .expect("the filtered-out tx must still resolve its block"); + assert_eq!(store.get_events_for_block(block_id).unwrap(), None); + } } #[cfg(test)] diff --git a/lez/indexer/core/src/event_filter.rs b/lez/indexer/core/src/event_filter.rs index ab688d6a0..2841b2951 100644 --- a/lez/indexer/core/src/event_filter.rs +++ b/lez/indexer/core/src/event_filter.rs @@ -34,6 +34,27 @@ impl EventFilter { } } + /// Whether every event in the requested `(program, selector)` domain is stored + /// under this filter; `None` widens the dimension to "all". + #[must_use] + pub fn covers(&self, program_id: Option, selector: Option<[u8; 8]>) -> bool { + match self { + Self::Archival => true, + Self::Sources(sources) => { + let Some(program_id) = program_id else { + return false; + }; + match sources.get(&program_id) { + None => false, + Some(SelectorFilter::All) => true, + Some(SelectorFilter::Only(selectors)) => { + selector.is_some_and(|selector| selectors.contains(&selector)) + } + } + } + } + } + #[must_use] pub fn filter_block(&self, block_events: Vec) -> Vec { if matches!(self, Self::Archival) { @@ -49,6 +70,36 @@ impl EventFilter { } } +/// Whether every event in the requested `(program, selector)` domain over +/// blocks `from..=to` was stored. +/// +/// Each filter segment whose span intersects the range must cover the domain, +/// and the range must not precede the first segment. +#[must_use] +pub fn covered_over_range( + segments: &[(EventFilter, u64)], + from: u64, + to: u64, + program_id: Option, + selector: Option<[u8; 8]>, +) -> bool { + if to < from { + return true; + } + let Some((_, first_from)) = segments.first() else { + return false; + }; + if from < *first_from { + return false; + } + segments.iter().enumerate().all(|(i, (filter, seg_from))| { + let seg_to = segments + .get(i.saturating_add(1)) + .map_or(u64::MAX, |(_, next_from)| next_from.saturating_sub(1)); + *seg_from > to || seg_to < from || filter.covers(program_id, selector) + }) +} + #[cfg(test)] mod tests { use common::HashType; @@ -151,4 +202,122 @@ mod tests { let expected = vec![group(1, vec![event(PROGRAM_A, SELECTOR_Y)])]; assert_eq!(filter.filter_block(blocks), expected); } + + #[test] + fn archival_covers_any_domain() { + for (program, selector) in [ + (None, None), + (Some(PROGRAM_A), None), + (Some(PROGRAM_A), Some(SELECTOR_X)), + ] { + assert!(EventFilter::Archival.covers(program, selector)); + } + } + + #[test] + fn default_covers_nothing() { + for (program, selector) in [ + (None, None), + (Some(PROGRAM_A), None), + (Some(PROGRAM_A), Some(SELECTOR_X)), + ] { + assert!(!EventFilter::default().covers(program, selector)); + } + } + + #[test] + fn undeclared_or_unspecified_program_is_not_covered() { + let filter = sources(vec![(PROGRAM_A, SelectorFilter::All)]); + + assert!(!filter.covers(Some(PROGRAM_B), Some(SELECTOR_X))); + assert!(!filter.covers(None, Some(SELECTOR_X))); + assert!(!filter.covers(None, None)); + } + + #[test] + fn program_wide_entry_covers_its_whole_domain() { + let filter = sources(vec![(PROGRAM_A, SelectorFilter::All)]); + + assert!(filter.covers(Some(PROGRAM_A), None)); + assert!(filter.covers(Some(PROGRAM_A), Some(SELECTOR_X))); + } + + #[test] + fn selector_entry_covers_only_listed_selectors() { + let filter = sources(vec![( + PROGRAM_A, + SelectorFilter::Only(HashSet::from([SELECTOR_X])), + )]); + + assert!(filter.covers(Some(PROGRAM_A), Some(SELECTOR_X))); + assert!(!filter.covers(Some(PROGRAM_A), Some(SELECTOR_Y))); + assert!(!filter.covers(Some(PROGRAM_A), None)); + } + + #[test] + fn range_coverage_needs_history() { + assert!(!covered_over_range(&[], 0, 10, None, None)); + } + + #[test] + fn range_before_the_first_segment_is_not_covered() { + let segments = [(EventFilter::Archival, 10)]; + + assert!(!covered_over_range(&segments, 5, 15, None, None)); + assert!(covered_over_range(&segments, 10, 15, None, None)); + } + + #[test] + fn every_segment_intersecting_the_range_must_cover_the_domain() { + let segments = [ + (sources(vec![(PROGRAM_A, SelectorFilter::All)]), 0), + ( + sources(vec![ + (PROGRAM_A, SelectorFilter::All), + (PROGRAM_B, SelectorFilter::All), + ]), + 100, + ), + ]; + + // Program B is only declared from block 100 on. + assert!(covered_over_range( + &segments, + 100, + 200, + Some(PROGRAM_B), + None + )); + assert!(!covered_over_range( + &segments, + 50, + 150, + Some(PROGRAM_B), + None + )); + assert!(!covered_over_range( + &segments, + 99, + 99, + Some(PROGRAM_B), + None + )); + // Program A is covered by both eras. + assert!(covered_over_range( + &segments, + 50, + 150, + Some(PROGRAM_A), + None + )); + // The whole-domain query needs archival coverage in every era. + assert!(!covered_over_range(&segments, 100, 200, None, None)); + } + + #[test] + fn an_empty_range_is_covered_vacuously() { + let segments = [(sources(vec![]), 0)]; + + assert!(covered_over_range(&segments, 5, 4, Some(PROGRAM_A), None)); + } } diff --git a/lez/indexer/service/src/service.rs b/lez/indexer/service/src/service.rs index 1f59f3562..135f7a6a1 100644 --- a/lez/indexer/service/src/service.rs +++ b/lez/indexer/service/src/service.rs @@ -7,10 +7,10 @@ use std::{ use anyhow::{Context as _, Result, bail}; use arc_swap::ArcSwap; use futures::StreamExt as _; -use indexer_core::{IndexerCore, config::IndexerConfig}; +use indexer_core::{IndexerCore, config::IndexerConfig, event_filter::EventFilter}; use indexer_service_protocol::{ Account, AccountId, Block, BlockId, EventRecord, EventSubscriptionFilter, GetEventsFilter, - HashType, IndexerStatus, Transaction, resolve_event_block_range, + HashType, IndexerStatus, ProgramId, Selector, Transaction, resolve_event_block_range, }; use jsonrpsee::{ SubscriptionSink, @@ -70,6 +70,14 @@ impl indexer_service_rpc::RpcServer for IndexerService { subscription_sink: jsonrpsee::PendingSubscriptionSink, filter: EventSubscriptionFilter, ) -> SubscriptionResult { + if let Err(err) = check_event_coverage( + self.indexer.store.live_filter(), + filter.program_id, + filter.selector, + ) { + subscription_sink.reject(err).await; + return Ok(()); + } let sink = subscription_sink.accept().await?; log::info!( "Accepted new subscription to events with ID {:?}", @@ -195,17 +203,46 @@ impl indexer_service_rpc::RpcServer for IndexerService { .unwrap_or(0); let records = match plan_query(&filter, tip)? { - EventQuery::ByTxHash(tx_hash) => self - .indexer - .store - .get_events_by_tx_hash(tx_hash.0) - .map_err(db_error)? - .map(|(block_id, group)| EventRecord::from_tx_events(block_id, group)) - .ok_or_else(unknown_transaction_error)?, + EventQuery::ByTxHash(tx_hash) => { + // Coverage is judged at the transaction's height, resolved BEFORE the + // events read: a filtered-out tx has no events row, and gating on the + // row's presence would serve `[]` for exactly the dropped domains. + let block_id = self + .indexer + .store + .block_id_by_tx_hash(tx_hash.0) + .map_err(db_error)? + .ok_or_else(unknown_transaction_error)?; + check_range_coverage( + self.indexer.store.filter_segments(), + block_id, + block_id, + filter.program_id, + filter.selector, + )?; + self.indexer + .store + .get_events_for_block(block_id) + .map_err(db_error)? + .and_then(|groups| { + groups + .into_iter() + .find(|group| group.tx_hash.0 == tx_hash.0) + }) + .map(|group| EventRecord::from_tx_events(block_id, group)) + .unwrap_or_default() + } EventQuery::ByRange { from_block, to_block, } => { + check_range_coverage( + self.indexer.store.filter_segments(), + from_block, + to_block, + filter.program_id, + filter.selector, + )?; let mut records = vec![]; for (block_id, groups) in self .indexer @@ -625,6 +662,53 @@ pub(crate) fn matches_subscription_filter( .is_none_or(|tx_hash| tx_hash == record.tx_hash) } +// Coverage is judged against the filters the store was WRITTEN under: a query +// they do not fully cover would be answered from a knowingly incomplete store. +// Subscriptions are forward-only, so they check the live filter. +pub(crate) fn check_event_coverage( + stored: &EventFilter, + program_id: Option, + selector: Option, +) -> Result<(), ErrorObjectOwned> { + if stored.covers(program_id.map(|id| id.0), selector.map(|s| s.0)) { + Ok(()) + } else { + Err(uncovered_query_error()) + } +} + +pub(crate) fn check_range_coverage( + segments: &[(EventFilter, u64)], + from: u64, + to: u64, + program_id: Option, + selector: Option, +) -> Result<(), ErrorObjectOwned> { + if indexer_core::event_filter::covered_over_range( + segments, + from, + to, + program_id.map(|id| id.0), + selector.map(|s| s.0), + ) { + Ok(()) + } else { + Err(uncovered_query_error()) + } +} + +fn uncovered_query_error() -> ErrorObjectOwned { + ErrorObjectOwned::owned( + ErrorCode::InvalidParams.code(), + "UncoveredEventQuery".to_owned(), + Some( + "this indexer's event filter does not cover the requested events; query a declared \ + program and selector, or an archival indexer" + .to_owned(), + ), + ) +} + fn invalid_params_error(message: impl Into) -> ErrorObjectOwned { ErrorObjectOwned::owned( ErrorCode::InvalidParams.code(), @@ -657,6 +741,9 @@ fn db_error(err: anyhow::Error) -> ErrorObjectOwned { #[cfg(test)] mod tests { + use std::collections::HashMap; + + use indexer_core::event_filter::SelectorFilter; use indexer_service_protocol::{MAX_EVENT_QUERY_BLOCK_SPAN, ProgramId, Selector}; use super::*; @@ -805,4 +892,33 @@ mod tests { assert!(!record(1, 7, 3).matches_fields(program, selector)); assert!(!record(1, 8, 2).matches_fields(program, selector)); } + + #[test] + fn coverage_check_accepts_archival_and_declared_sources() { + assert!(check_event_coverage(&EventFilter::Archival, None, None).is_ok()); + + let declared = EventFilter::Sources(HashMap::from([([7; 8], SelectorFilter::All)])); + assert!( + check_event_coverage(&declared, Some(ProgramId([7; 8])), Some(Selector([1; 8]))) + .is_ok() + ); + } + + #[test] + fn range_coverage_follows_segment_history() { + let declared = EventFilter::Sources(HashMap::from([([7; 8], SelectorFilter::All)])); + let segments = [(declared, 0), (EventFilter::Archival, 100)]; + + assert!(check_range_coverage(&segments, 100, 200, None, None).is_ok()); + assert!(check_range_coverage(&segments, 50, 150, Some(ProgramId([7; 8])), None).is_ok()); + let err = check_range_coverage(&segments, 50, 150, None, None).unwrap_err(); + assert_eq!(err.message(), "UncoveredEventQuery"); + } + + #[test] + fn uncovered_query_is_rejected() { + let err = check_event_coverage(&EventFilter::default(), Some(ProgramId([7; 8])), None) + .unwrap_err(); + assert_eq!(err.message(), "UncoveredEventQuery"); + } }