From bb47bdaa4efb377f2a2deb775ef4d6b16f58206f Mon Sep 17 00:00:00 2001 From: Jerry Date: Fri, 13 Sep 2024 10:54:47 -0700 Subject: [PATCH] Add pending logs back --- eth/filters/filter_system.go | 17 ++++++++++++++++- 1 file changed, 16 insertions(+), 1 deletion(-) diff --git a/eth/filters/filter_system.go b/eth/filters/filter_system.go index ed6153989e..5fb7b2cbe8 100644 --- a/eth/filters/filter_system.go +++ b/eth/filters/filter_system.go @@ -188,7 +188,7 @@ const ( // logsChanSize is the size of channel listening to LogsEvent. // Updated to fix TestEth2NeBlock testcase, as the feed was unable to send // logs to the channel. todo: @anshalshukla - logsChanSize = 100 + logsChanSize = 10 // chainEvChanSize is the size of channel listening to ChainEvent. chainEvChanSize = 10 // stateEvChanSize is the size of channel listening to StateSyncEvent. @@ -427,6 +427,19 @@ func (es *EventSystem) handleLogs(filters filterIndex, ev []*types.Log) { } } +func (es *EventSystem) handlePendingLogs(filters filterIndex, ev []*types.Log) { + if len(ev) == 0 { + return + } + + for _, f := range filters[PendingLogsSubscription] { + matchedLogs := filterLogs(ev, nil, f.logsCrit.ToBlock, f.logsCrit.Addresses, f.logsCrit.Topics) + if len(matchedLogs) > 0 { + f.logs <- matchedLogs + } + } +} + func (es *EventSystem) handleTxsEvent(filters filterIndex, ev core.NewTxsEvent) { for _, f := range filters[PendingTransactionsSubscription] { f.txs <- ev.Txs @@ -464,6 +477,8 @@ func (es *EventSystem) eventLoop() { es.handleLogs(index, ev) case ev := <-es.rmLogsCh: es.handleLogs(index, ev.Logs) + case ev := <-es.pendingLogsCh: + es.handlePendingLogs(index, ev) case ev := <-es.chainCh: es.handleChainEvent(index, ev) case ev := <-es.stateSyncCh: