@@ -52,9 +52,15 @@ use fvm_ipld_encoding::IPLD_RAW;
5252use serde:: * ;
5353use std:: ops:: RangeInclusive ;
5454use std:: sync:: Arc ;
55+ use std:: time:: Duration ;
5556use store:: * ;
5657use tokio_util:: task:: AbortOnDropHandle ;
5758
59+ /// Longest period between expired-filter sweeps; a TTL shorter than this is
60+ /// swept every TTL instead. Mirrors Lotus's hardcoded 30-minute GC ticker:
61+ /// <https://github.com/filecoin-project/lotus/blob/release/v1.36.1/node/impl/eth/events.go#L361>
62+ const FILTER_GC_INTERVAL : Duration = Duration :: from_mins ( 30 ) ;
63+
5864/// A trait for filtering events based on predefined conditions.
5965///
6066/// Implementors of this trait define custom logic to determine whether an event matches the filtering criteria
@@ -120,6 +126,7 @@ pub struct EthEventHandler {
120126 pub filter_store : Option < Arc < dyn FilterStore > > ,
121127 pub max_filter_results : usize ,
122128 pub max_filter_height_range : ChainEpoch ,
129+ filter_ttl : Option < Duration > ,
123130 event_filter_manager : Option < Arc < EventFilterManager > > ,
124131 tipset_filter_manager : Option < Arc < TipSetFilterManager > > ,
125132 mempool_filter_manager : Option < Arc < MempoolFilterManager > > ,
@@ -173,6 +180,8 @@ impl EthEventHandler {
173180 }
174181 } )
175182 . unwrap_or ( config. max_filter_height_range ) ;
183+ let filter_ttl_secs: u64 = env_or_default ( "FOREST_FILTER_TTL_SECS" , config. filter_ttl_secs ) ;
184+ let filter_ttl = ( filter_ttl_secs > 0 ) . then ( || Duration :: from_secs ( filter_ttl_secs) ) ;
176185 let filter_store: Option < Arc < dyn FilterStore > > =
177186 Some ( MemFilterStore :: new ( max_filters) as Arc < dyn FilterStore > ) ;
178187 let event_filter_manager = Some ( EventFilterManager :: new ( ) ) ;
@@ -187,12 +196,62 @@ impl EthEventHandler {
187196 filter_store,
188197 max_filter_results,
189198 max_filter_height_range,
199+ filter_ttl,
190200 event_filter_manager,
191201 tipset_filter_manager,
192202 mempool_filter_manager,
193203 }
194204 }
195205
206+ /// TTL after which idle filters are removed; `None` when expiry is disabled.
207+ pub fn filter_ttl ( & self ) -> Option < Duration > {
208+ self . filter_ttl
209+ }
210+
211+ /// Periodically removes filters idle for longer than `ttl`, sweeping every
212+ /// [`FILTER_GC_INTERVAL`] or every `ttl`, whichever is shorter.
213+ pub async fn run_filter_gc ( & self , ttl : Duration ) -> ! {
214+ let period = ttl. min ( FILTER_GC_INTERVAL ) ;
215+ let mut interval = tokio:: time:: interval_at ( tokio:: time:: Instant :: now ( ) + period, period) ;
216+ loop {
217+ interval. tick ( ) . await ;
218+ self . sweep_expired ( ttl) ;
219+ }
220+ }
221+
222+ /// One sweep pass: removes every filter idle for longer than `ttl` from
223+ /// the store and from its manager.
224+ fn sweep_expired ( & self , ttl : Duration ) {
225+ let Some ( filter_store) = & self . filter_store else {
226+ return ;
227+ } ;
228+ let expired = filter_store. remove_expired ( ttl) ;
229+ for filter in & expired {
230+ if let Err ( e) = self . remove_from_manager ( filter. as_ref ( ) ) {
231+ tracing:: warn!(
232+ "Failed to remove expired eth filter {:?} from its manager: {e:#}" ,
233+ filter. id( )
234+ ) ;
235+ }
236+ }
237+ if !expired. is_empty ( ) {
238+ tracing:: debug!( "Expired {} idle eth filter(s)" , expired. len( ) ) ;
239+ }
240+ }
241+
242+ /// Completes an install: adds `filter` to the store and returns its id,
243+ /// rolling the manager registration back when the store rejects it (e.g.
244+ /// at the filter cap).
245+ fn register_filter ( & self , filter : Arc < dyn Filter > ) -> Result < FilterID , Error > {
246+ if let Some ( filter_store) = & self . filter_store
247+ && let Err ( err) = filter_store. add ( filter. clone ( ) )
248+ {
249+ self . remove_from_manager ( filter. as_ref ( ) ) ?;
250+ bail ! ( "Adding filter failed: {}" , err) ;
251+ }
252+ Ok ( filter. id ( ) . clone ( ) )
253+ }
254+
196255 // Installs an eth filter based on given filter spec.
197256 pub fn eth_new_filter (
198257 & self ,
@@ -208,16 +267,7 @@ impl EthEventHandler {
208267 . install ( pf)
209268 . context ( "Installation error" ) ?;
210269
211- if let Some ( filter_store) = & self . filter_store
212- && let Err ( err) = filter_store. add ( filter. clone ( ) )
213- {
214- ensure ! (
215- event_filter_manager. remove( filter. id( ) ) . is_some( ) ,
216- "Filter not found"
217- ) ;
218- bail ! ( "Adding filter failed: {}" , err) ;
219- }
220- Ok ( filter. id ( ) . clone ( ) )
270+ self . register_filter ( filter)
221271 } else {
222272 Err ( Error :: msg ( "NotSupported" ) )
223273 }
@@ -229,13 +279,7 @@ impl EthEventHandler {
229279 ) -> Result < FilterID , Error > {
230280 if let Some ( manager) = filter_manager {
231281 let filter = manager. install ( ) . context ( "Installation error" ) ?;
232- if let Some ( filter_store) = & self . filter_store
233- && let Err ( err) = filter_store. add ( filter. clone ( ) )
234- {
235- ensure ! ( manager. remove( filter. id( ) ) . is_some( ) , "Filter not found" ) ;
236- bail ! ( "Adding filter failed: {}" , err) ;
237- }
238- Ok ( filter. id ( ) . clone ( ) )
282+ self . register_filter ( filter)
239283 } else {
240284 Err ( Error :: msg ( "NotSupported" ) )
241285 }
@@ -251,7 +295,7 @@ impl EthEventHandler {
251295 self . install_filter ( self . mempool_filter_manager . as_deref ( ) . map ( |fm| fm as _ ) )
252296 }
253297
254- fn uninstall_filter ( & self , filter : Arc < dyn Filter > ) -> Result < ( ) , Error > {
298+ fn remove_from_manager ( & self , filter : & dyn Filter ) -> Result < ( ) , Error > {
255299 let id = filter. id ( ) ;
256300
257301 if filter. as_any ( ) . is :: < EventFilter > ( ) {
@@ -274,23 +318,18 @@ impl EthEventHandler {
274318 . context ( "Failed to remove mempool filter" ) ?;
275319 }
276320
277- self . filter_store
278- . as_ref ( )
279- . context ( "Filter store is missing" ) ?
280- . remove ( id)
281- . context ( "Failed to remove filter from store" ) ?;
282-
283321 Ok ( ( ) )
284322 }
285323
324+ /// Uninstalls an eth filter.
286325 pub fn eth_uninstall_filter ( & self , id : & FilterID ) -> Result < bool , Error > {
287326 let store = self
288327 . filter_store
289328 . as_ref ( )
290329 . context ( "Filter store is not supported" ) ?;
291330
292- if let Ok ( filter) = store. get ( id) {
293- self . uninstall_filter ( filter) ?;
331+ if let Some ( filter) = store. remove ( id) {
332+ self . remove_from_manager ( filter. as_ref ( ) ) ?;
294333 Ok ( true )
295334 } else {
296335 Ok ( false )
@@ -902,8 +941,8 @@ mod tests {
902941 assert ! ( parsed. addresses. is_empty( ) ) ;
903942 }
904943
905- #[ test]
906- fn test_eth_new_filter_with_none_address ( ) {
944+ #[ tokio :: test]
945+ async fn test_eth_new_filter_with_none_address ( ) {
907946 let eth_event_handler = EthEventHandler :: new ( ) ;
908947
909948 let filter_spec = EthFilterSpec {
@@ -1209,8 +1248,8 @@ mod tests {
12091248 assert ! ( result. is_err( ) ) ;
12101249 }
12111250
1212- #[ test]
1213- fn test_eth_new_filter ( ) {
1251+ #[ tokio :: test]
1252+ async fn test_eth_new_filter ( ) {
12141253 let eth_event_handler = EthEventHandler :: new ( ) ;
12151254
12161255 let filter_spec = EthFilterSpec {
@@ -1229,8 +1268,8 @@ mod tests {
12291268 assert ! ( result. is_ok( ) , "Expected successful filter creation" ) ;
12301269 }
12311270
1232- #[ test]
1233- fn test_eth_new_block_filter ( ) {
1271+ #[ tokio :: test]
1272+ async fn test_eth_new_block_filter ( ) {
12341273 let eth_event_handler = EthEventHandler :: new ( ) ;
12351274 let result = eth_event_handler. eth_new_block_filter ( ) ;
12361275
@@ -1269,10 +1308,121 @@ mod tests {
12691308 let pending_tx_filter_id = event_handler. eth_new_pending_transaction_filter ( ) . unwrap ( ) ;
12701309 filter_ids. push ( pending_tx_filter_id) ;
12711310
1272- for filter_id in filter_ids {
1273- let result = event_handler. eth_uninstall_filter ( & filter_id) . unwrap ( ) ;
1311+ for filter_id in & filter_ids {
1312+ let result = event_handler. eth_uninstall_filter ( filter_id) . unwrap ( ) ;
12741313 assert ! ( result, "Uninstalling filter with id {filter_id:?} failed" ) ;
12751314 }
1315+
1316+ // uninstalling an already-removed filter reports `false` instead of erroring
1317+ let gone = filter_ids. first ( ) . unwrap ( ) ;
1318+ assert ! ( !event_handler. eth_uninstall_filter( gone) . unwrap( ) ) ;
1319+ }
1320+
1321+ const TEST_TTL : Duration = Duration :: from_hours ( 1 ) ;
1322+
1323+ /// Handler whose filters expire after `ttl` without being polled;
1324+ /// `Duration::ZERO` disables expiry.
1325+ fn handler_with_ttl ( ttl : Duration ) -> EthEventHandler {
1326+ EthEventHandler :: from_config (
1327+ & EventsConfig {
1328+ filter_ttl_secs : ttl. as_secs ( ) ,
1329+ ..Default :: default ( )
1330+ } ,
1331+ crate :: networks:: mainnet:: ETH_CHAIN_ID ,
1332+ MpoolSubscriber :: dummy ( ) ,
1333+ )
1334+ }
1335+
1336+ #[ tokio:: test( start_paused = true ) ]
1337+ async fn sweep_expired_removes_idle_filters_of_every_type ( ) {
1338+ let handler = handler_with_ttl ( TEST_TTL ) ;
1339+ let event_id = handler
1340+ . eth_new_filter ( & EthFilterSpec :: default ( ) , 0 )
1341+ . unwrap ( ) ;
1342+ let block_id = handler. eth_new_block_filter ( ) . unwrap ( ) ;
1343+ let pending_id = handler. eth_new_pending_transaction_filter ( ) . unwrap ( ) ;
1344+
1345+ // everything has been idle for longer than the TTL
1346+ tokio:: time:: advance ( TEST_TTL + Duration :: from_secs ( 1 ) ) . await ;
1347+ handler. sweep_expired ( TEST_TTL ) ;
1348+
1349+ // the filters are gone from the store...
1350+ let store = handler. filter_store . as_ref ( ) . unwrap ( ) ;
1351+ for id in [ & event_id, & block_id, & pending_id] {
1352+ assert ! ( store. get( id) . is_err( ) , "filter {id:?} should be expired" ) ;
1353+ }
1354+ // ...and from their managers
1355+ let event_manager = handler. event_filter_manager . as_ref ( ) . unwrap ( ) ;
1356+ assert ! ( event_manager. remove( & event_id) . is_none( ) ) ;
1357+ let tipset_manager = handler. tipset_filter_manager . as_ref ( ) . unwrap ( ) ;
1358+ assert ! ( tipset_manager. remove( & block_id) . is_none( ) ) ;
1359+ let mempool_manager = handler. mempool_filter_manager . as_ref ( ) . unwrap ( ) ;
1360+ assert ! ( mempool_manager. remove( & pending_id) . is_none( ) ) ;
1361+ }
1362+
1363+ #[ tokio:: test( start_paused = true ) ]
1364+ async fn sweep_expired_spares_recently_polled_filters ( ) {
1365+ let handler = handler_with_ttl ( TEST_TTL ) ;
1366+ let polled_id = handler. eth_new_block_filter ( ) . unwrap ( ) ;
1367+ let idle_id = handler. eth_new_block_filter ( ) . unwrap ( ) ;
1368+ let store = handler. filter_store . as_ref ( ) . unwrap ( ) ;
1369+
1370+ // poll one filter 45 minutes in
1371+ tokio:: time:: advance ( Duration :: from_mins ( 45 ) ) . await ;
1372+ store. get ( & polled_id) . unwrap ( ) ;
1373+
1374+ // 30 minutes later, the never-polled filter has crossed the TTL
1375+ // (75 minutes idle) while the polled one has not (30 minutes idle)
1376+ tokio:: time:: advance ( Duration :: from_mins ( 30 ) ) . await ;
1377+ handler. sweep_expired ( TEST_TTL ) ;
1378+
1379+ assert ! ( store. get( & idle_id) . is_err( ) , "idle filter should expire" ) ;
1380+ assert ! (
1381+ store. get( & polled_id) . is_ok( ) ,
1382+ "recently polled filter must survive"
1383+ ) ;
1384+ }
1385+
1386+ #[ tokio:: test( start_paused = true ) ]
1387+ async fn gc_loop_sweeps_idle_filters ( ) {
1388+ let handler = Arc :: new ( handler_with_ttl ( TEST_TTL ) ) ;
1389+ let id = handler. eth_new_block_filter ( ) . unwrap ( ) ;
1390+ let _gc = tokio:: spawn ( {
1391+ let handler = handler. clone ( ) ;
1392+ async move { handler. run_filter_gc ( TEST_TTL ) . await }
1393+ } ) ;
1394+ tokio:: task:: yield_now ( ) . await ; // let the spawned loop start its interval now
1395+
1396+ tokio:: time:: advance ( TEST_TTL + FILTER_GC_INTERVAL ) . await ;
1397+ tokio:: task:: yield_now ( ) . await ; // let the loop process elapsed ticks
1398+ let store = handler. filter_store . as_ref ( ) . unwrap ( ) ;
1399+ assert ! ( store. get( & id) . is_err( ) , "idle filter should be swept" ) ;
1400+ }
1401+
1402+ #[ tokio:: test( start_paused = true ) ]
1403+ async fn short_ttl_sweeps_before_the_default_interval ( ) {
1404+ let ttl = Duration :: from_mins ( 5 ) ;
1405+ let handler = Arc :: new ( handler_with_ttl ( ttl) ) ;
1406+ let id = handler. eth_new_block_filter ( ) . unwrap ( ) ;
1407+ let _gc = tokio:: spawn ( {
1408+ let handler = handler. clone ( ) ;
1409+ async move { handler. run_filter_gc ( ttl) . await }
1410+ } ) ;
1411+ tokio:: task:: yield_now ( ) . await ; // let the spawned loop start its interval now
1412+
1413+ // two TTL periods, well inside the 30-minute `FILTER_GC_INTERVAL`
1414+ tokio:: time:: advance ( ttl * 2 + Duration :: from_secs ( 1 ) ) . await ;
1415+ tokio:: task:: yield_now ( ) . await ; // let the loop process elapsed ticks
1416+ let store = handler. filter_store . as_ref ( ) . unwrap ( ) ;
1417+ assert ! (
1418+ store. get( & id) . is_err( ) ,
1419+ "short-TTL filter must be swept before the default interval elapses"
1420+ ) ;
1421+ }
1422+
1423+ #[ test]
1424+ fn zero_ttl_disables_filter_expiry ( ) {
1425+ assert ! ( handler_with_ttl( Duration :: ZERO ) . filter_ttl( ) . is_none( ) ) ;
12761426 }
12771427
12781428 #[ test]
0 commit comments