feat(indexer): reject event queries and subscriptions outside the declared filter

This commit is contained in:
Artem Gureev
2026-08-23 06:01:42 +00:00
parent 6b469e9690
commit f4d3e96632
3 changed files with 356 additions and 25 deletions
+62 -16
View File
@@ -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<Option<(u64, TxEvents)>> {
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<Option<u64>> {
Ok(self.dbio.get_block_id_by_tx_hash(tx_hash)?)
}
pub fn get_block_by_hash(&self, hash: [u8; 32]) -> Result<Option<Block>> {
@@ -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)]
+169
View File
@@ -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<ProgramId>, 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<TxEvents>) -> Vec<TxEvents> {
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<ProgramId>,
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));
}
}
+125 -9
View File
@@ -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<ProgramId>,
selector: Option<Selector>,
) -> 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<ProgramId>,
selector: Option<Selector>,
) -> 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<String>) -> 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");
}
}