diff --git a/src/rpc/methods/eth.rs b/src/rpc/methods/eth.rs
index dfe0cddeb0e..e6a597c489f 100644
--- a/src/rpc/methods/eth.rs
+++ b/src/rpc/methods/eth.rs
@@ -3447,37 +3447,7 @@ async fn poll_event_filter(
SkipEvent::OnUnresolvedAddress,
)
.await?;
- let mut seen_positions = SeenEventPositions::default();
- let mut recent_events = Vec::new();
- for event in events {
- let position = (event.msg_idx, event.event_idx);
- let already_seen = event_filter
- .seen_positions
- .get(&event.tipset_key)
- .is_some_and(|positions| positions.contains(&position));
- match seen_positions.get_mut(&event.tipset_key) {
- Some(positions) => {
- positions.insert(position);
- }
- None => {
- seen_positions.insert(event.tipset_key.clone(), HashSet::from_iter([position]));
- }
- }
- if !already_seen {
- recent_events.push(event);
- }
- }
- if let Some(store) = &ctx.eth_event_handler.filter_store {
- store.update(Arc::new(EventFilter {
- id: event_filter.id.clone(),
- tipsets: event_filter.tipsets.clone(),
- addresses: event_filter.addresses.clone(),
- keys_with_codec: event_filter.keys_with_codec.clone(),
- max_results: event_filter.max_results,
- seen_positions,
- }));
- }
- Ok(recent_events)
+ Ok(event_filter.take_unseen(events))
}
pub enum EthGetFilterLogs {}
@@ -3541,7 +3511,7 @@ impl RpcMethod<1> for EthGetFilterChanges {
// heaviest tipset doesn't have events because its messages haven't been executed yet
RangeInclusive::new(
tipset_filter
- .collected
+ .collected()
.unwrap_or(ctx.chain_store().heaviest_tipset().epoch() - 1),
// Use -1 to indicate that the range extends until the latest available tipset.
-1,
@@ -3555,12 +3525,7 @@ impl RpcMethod<1> for EthGetFilterChanges {
.max_by_key(|event| event.height)
.map(|e| e.height);
if let Some(height) = new_collected {
- let filter = Arc::new(TipSetFilter {
- id: tipset_filter.id.clone(),
- max_results: tipset_filter.max_results,
- collected: Some(height),
- });
- store.update(filter);
+ tipset_filter.set_collected(height);
}
return Ok(eth_filter_result_from_tipsets(&events)?);
}
diff --git a/src/rpc/methods/eth/filter/event.rs b/src/rpc/methods/eth/filter/event.rs
index 14b871bd358..9d6b0633356 100644
--- a/src/rpc/methods/eth/filter/event.rs
+++ b/src/rpc/methods/eth/filter/event.rs
@@ -3,14 +3,14 @@
use crate::prelude::*;
use crate::rpc::eth::filter::{ActorEventBlock, ParsedFilter, ParsedFilterTipsets};
-use crate::rpc::eth::{FilterID, SeenEventPositions, filter::Filter};
+use crate::rpc::eth::{CollectedEvent, FilterID, SeenEventPositions, filter::Filter};
use crate::shim::address::Address;
-use ahash::HashMap;
+use ahash::{HashMap, HashSet};
use anyhow::Result;
-use parking_lot::RwLock;
+use parking_lot::{Mutex, RwLock};
use std::any::Any;
-#[derive(Debug, PartialEq)]
+#[derive(Debug)]
pub struct EventFilter {
// Unique id used to identify the filter
pub id: FilterID,
@@ -20,10 +20,37 @@ pub struct EventFilter {
pub addresses: Vec
,
// Map of key names to a list of alternate values that may match
pub keys_with_codec: HashMap>,
- // Maximum number of results to collect
- pub max_results: usize,
// Positions of the events returned by the last poll, used to compute the next poll's delta
- pub seen_positions: SeenEventPositions,
+ seen_positions: Mutex,
+}
+
+impl EventFilter {
+ /// Records the positions of `events` as the filter's new poll cursor and
+ /// returns only the events the previous poll did not contain.
+ pub fn take_unseen(&self, events: Vec) -> Vec {
+ let mut seen_positions = self.seen_positions.lock();
+ let mut new_positions = SeenEventPositions::default();
+ let mut recent_events = Vec::new();
+ for event in events {
+ let position = (event.msg_idx, event.event_idx);
+ let already_seen = seen_positions
+ .get(&event.tipset_key)
+ .is_some_and(|positions| positions.contains(&position));
+ match new_positions.get_mut(&event.tipset_key) {
+ Some(positions) => {
+ positions.insert(position);
+ }
+ None => {
+ new_positions.insert(event.tipset_key.clone(), HashSet::from_iter([position]));
+ }
+ }
+ if !already_seen {
+ recent_events.push(event);
+ }
+ }
+ *seen_positions = new_positions;
+ recent_events
+ }
}
impl From<&EventFilter> for ParsedFilter {
@@ -49,17 +76,15 @@ impl Filter for EventFilter {
/// The `EventFilterManager` structure maintains a set of filters, allowing new filters to be
/// installed or existing ones to be removed. It ensures that each filter is uniquely identifiable
-/// by its ID and that a maximum number of results can be configured for each filter.
+/// by its ID.
pub struct EventFilterManager {
filters: RwLock>>,
- max_filter_results: usize,
}
impl EventFilterManager {
- pub fn new(max_filter_results: usize) -> Arc {
+ pub fn new() -> Arc {
Arc::new(Self {
filters: RwLock::new(HashMap::new()),
- max_filter_results,
})
}
@@ -71,7 +96,6 @@ impl EventFilterManager {
tipsets: pf.tipsets,
addresses: pf.addresses,
keys_with_codec: pf.keys,
- max_results: self.max_filter_results,
seen_positions: Default::default(),
});
@@ -89,14 +113,82 @@ impl EventFilterManager {
#[cfg(test)]
mod tests {
use super::*;
+ use crate::blocks::TipsetKey;
use crate::rpc::eth::filter::{ParsedFilter, ParsedFilterTipsets};
use crate::shim::address::Address;
+ use crate::utils::multihash::MultihashCode;
+ use fvm_ipld_encoding::DAG_CBOR;
+ use multihash_derive::MultihashDigest as _;
+ use nunny::vec as nonempty;
use std::ops::RangeInclusive;
+ fn install_filter() -> Arc {
+ EventFilterManager::new()
+ .install(ParsedFilter {
+ tipsets: ParsedFilterTipsets::Range(RangeInclusive::new(0, 100)),
+ addresses: vec![],
+ keys: HashMap::new(),
+ msg_cid: None,
+ })
+ .expect("Failed to install EventFilter")
+ }
+
+ fn tipset_key(seed: u8) -> TipsetKey {
+ TipsetKey::from(nonempty![Cid::new_v1(
+ DAG_CBOR,
+ MultihashCode::Identity.digest(&[seed])
+ )])
+ }
+
+ fn event_at(tipset_key: &TipsetKey, msg_idx: u64, event_idx: u64) -> CollectedEvent {
+ CollectedEvent {
+ entries: vec![],
+ emitter_addr: Address::new_id(0),
+ event_idx,
+ reverted: false,
+ height: 0,
+ tipset_key: tipset_key.clone(),
+ msg_idx,
+ msg_cid: Cid::new_v1(DAG_CBOR, MultihashCode::Identity.digest(&[])),
+ }
+ }
+
+ #[test]
+ fn take_unseen_returns_only_new_events() {
+ let filter = install_filter();
+ let ts = tipset_key(1);
+ let (e0, e1) = (event_at(&ts, 0, 0), event_at(&ts, 0, 1));
+
+ // first poll: everything is new
+ assert_eq!(
+ filter.take_unseen(vec![e0.clone(), e1.clone()]),
+ vec![e0.clone(), e1.clone()]
+ );
+ // same result set again, nothing new
+ assert!(filter.take_unseen(vec![e0.clone(), e1.clone()]).is_empty());
+ // a new event alongside the old ones: only the new one comes back
+ let e2 = event_at(&ts, 1, 0);
+ assert_eq!(filter.take_unseen(vec![e0, e1, e2.clone()]), vec![e2]);
+ }
+
+ #[test]
+ fn take_unseen_cursor_tracks_only_the_last_poll() {
+ let filter = install_filter();
+ let (a, b) = (
+ event_at(&tipset_key(1), 0, 0),
+ event_at(&tipset_key(2), 0, 0),
+ );
+
+ assert_eq!(filter.take_unseen(vec![a.clone()]), vec![a.clone()]);
+ // a poll without tipset 1 in its results forgets its positions...
+ assert_eq!(filter.take_unseen(vec![b.clone()]), vec![b.clone()]);
+ // ...so its event counts as new when it reappears, while tipset 2's does not
+ assert_eq!(filter.take_unseen(vec![a.clone(), b]), vec![a]);
+ }
+
#[test]
fn test_event_filter() {
- let max_filter_results = 10;
- let event_manager = EventFilterManager::new(max_filter_results);
+ let event_manager = EventFilterManager::new();
let parsed_filter = ParsedFilter {
tipsets: ParsedFilterTipsets::Range(RangeInclusive::new(0, 100)),
@@ -119,8 +211,8 @@ mod tests {
// Test case 2: Remove the EventFilter
let removed = event_manager.remove(&filter_id);
assert_eq!(
- removed,
- Some(filter),
+ removed.map(|f| f.id().clone()),
+ Some(filter_id.clone()),
"Filter should be successfully removed"
);
diff --git a/src/rpc/methods/eth/filter/mod.rs b/src/rpc/methods/eth/filter/mod.rs
index 23dc2a3fb1f..ae061fc5fc6 100644
--- a/src/rpc/methods/eth/filter/mod.rs
+++ b/src/rpc/methods/eth/filter/mod.rs
@@ -175,8 +175,8 @@ impl EthEventHandler {
.unwrap_or(config.max_filter_height_range);
let filter_store: Option> =
Some(MemFilterStore::new(max_filters) as Arc);
- let event_filter_manager = Some(EventFilterManager::new(max_filter_results));
- let tipset_filter_manager = Some(TipSetFilterManager::new(max_filter_results));
+ let event_filter_manager = Some(EventFilterManager::new());
+ let tipset_filter_manager = Some(TipSetFilterManager::new());
let mempool_filter_manager = Some(MempoolFilterManager::new(
max_filter_results,
eth_chain_id,
diff --git a/src/rpc/methods/eth/filter/store.rs b/src/rpc/methods/eth/filter/store.rs
index fca935bc9f8..f5534ce64c3 100644
--- a/src/rpc/methods/eth/filter/store.rs
+++ b/src/rpc/methods/eth/filter/store.rs
@@ -21,7 +21,6 @@ pub trait Filter: Send + Sync + std::fmt::Debug {
pub trait FilterStore: Send + Sync {
fn add(&self, filter: Arc) -> Result<()>;
fn get(&self, id: &FilterID) -> Result>;
- fn update(&self, filter: Arc);
fn remove(&self, id: &FilterID) -> Option>;
}
@@ -62,12 +61,6 @@ impl FilterStore for MemFilterStore {
.ok_or_else(|| anyhow!("filter not found"))
}
- fn update(&self, filter: Arc) {
- let mut filters = self.filters.write();
-
- filters.insert(filter.id().clone(), filter);
- }
-
fn remove(&self, id: &FilterID) -> Option> {
let mut filters = self.filters.write();
filters.remove(id)
diff --git a/src/rpc/methods/eth/filter/tipset.rs b/src/rpc/methods/eth/filter/tipset.rs
index 4a2e5d51c2a..3423be645dc 100644
--- a/src/rpc/methods/eth/filter/tipset.rs
+++ b/src/rpc/methods/eth/filter/tipset.rs
@@ -6,29 +6,35 @@ use crate::rpc::eth::{FilterID, filter::Filter, filter::FilterManager};
use crate::shim::fvm_shared_latest::clock::ChainEpoch;
use ahash::HashMap;
use anyhow::Result;
-use parking_lot::RwLock;
+use parking_lot::{Mutex, RwLock};
use std::any::Any;
-#[allow(dead_code)]
-#[derive(Debug, PartialEq)]
+#[derive(Debug)]
pub struct TipSetFilter {
// Unique id used to identify the filter
pub id: FilterID,
- // Maximum number of results to collect
- pub max_results: usize,
// Epoch at which the results were collected
- pub collected: Option,
+ collected: Mutex