mirror of
https://gitlab.com/pulsechaincom/erigon-pulse.git
synced 2024-12-31 16:21:21 +00:00
rpc: fix for map concurrency issue in logs subscription (#4903)
moving a couple of mutex locks and introducing another to prevent a deferred call to unsubscribe clashing with a new call to subscribe
This commit is contained in:
parent
ec67e80a8a
commit
98c639784b
@ -17,12 +17,13 @@ import (
|
||||
"github.com/ledgerwatch/erigon-lib/gointerfaces/remote"
|
||||
"github.com/ledgerwatch/erigon-lib/gointerfaces/txpool"
|
||||
txpool2 "github.com/ledgerwatch/erigon-lib/txpool"
|
||||
"github.com/ledgerwatch/log/v3"
|
||||
"google.golang.org/grpc"
|
||||
|
||||
"github.com/ledgerwatch/erigon/common"
|
||||
"github.com/ledgerwatch/erigon/core/types"
|
||||
"github.com/ledgerwatch/erigon/eth/filters"
|
||||
"github.com/ledgerwatch/erigon/rlp"
|
||||
"github.com/ledgerwatch/log/v3"
|
||||
"google.golang.org/grpc"
|
||||
)
|
||||
|
||||
type (
|
||||
@ -400,14 +401,14 @@ func (ff *Filters) SubscribeLogs(out chan *types.Log, crit filters.FilterCriteri
|
||||
AllAddresses: ff.logsSubs.aggLogsFilter.allAddrs == 1,
|
||||
AllTopics: ff.logsSubs.aggLogsFilter.allTopics == 1,
|
||||
}
|
||||
ff.mu.Lock()
|
||||
defer ff.mu.Unlock()
|
||||
for addr := range ff.logsSubs.aggLogsFilter.addrs {
|
||||
lfr.Addresses = append(lfr.Addresses, gointerfaces.ConvertAddressToH160(addr))
|
||||
}
|
||||
for topic := range ff.logsSubs.aggLogsFilter.topics {
|
||||
lfr.Topics = append(lfr.Topics, gointerfaces.ConvertHashToH256(topic))
|
||||
}
|
||||
ff.mu.Lock()
|
||||
defer ff.mu.Unlock()
|
||||
loaded := ff.logsRequestor.Load()
|
||||
if loaded != nil {
|
||||
if err := loaded.(func(*remote.LogsFilterRequest) error)(lfr); err != nil {
|
||||
@ -424,14 +425,14 @@ func (ff *Filters) UnsubscribeLogs(id LogsSubID) bool {
|
||||
AllAddresses: ff.logsSubs.aggLogsFilter.allAddrs == 1,
|
||||
AllTopics: ff.logsSubs.aggLogsFilter.allTopics == 1,
|
||||
}
|
||||
ff.mu.Lock()
|
||||
defer ff.mu.Unlock()
|
||||
for addr := range ff.logsSubs.aggLogsFilter.addrs {
|
||||
lfr.Addresses = append(lfr.Addresses, gointerfaces.ConvertAddressToH160(addr))
|
||||
}
|
||||
for topic := range ff.logsSubs.aggLogsFilter.topics {
|
||||
lfr.Topics = append(lfr.Topics, gointerfaces.ConvertHashToH256(topic))
|
||||
}
|
||||
ff.mu.Lock()
|
||||
defer ff.mu.Unlock()
|
||||
loaded := ff.logsRequestor.Load()
|
||||
if loaded != nil {
|
||||
if err := loaded.(func(*remote.LogsFilterRequest) error)(lfr); err != nil {
|
||||
|
@ -5,6 +5,7 @@ import (
|
||||
|
||||
"github.com/ledgerwatch/erigon-lib/gointerfaces"
|
||||
"github.com/ledgerwatch/erigon-lib/gointerfaces/remote"
|
||||
|
||||
"github.com/ledgerwatch/erigon/common"
|
||||
types2 "github.com/ledgerwatch/erigon/core/types"
|
||||
)
|
||||
@ -81,6 +82,8 @@ func (a *LogsFilterAggregator) subtractLogFilters(f *LogsFilter) {
|
||||
}
|
||||
|
||||
func (a *LogsFilterAggregator) addLogsFilters(f *LogsFilter) {
|
||||
a.logsFilterLock.Lock()
|
||||
defer a.logsFilterLock.Unlock()
|
||||
a.aggLogsFilter.allAddrs += f.allAddrs
|
||||
for addr, count := range f.addrs {
|
||||
a.aggLogsFilter.addrs[addr] += count
|
||||
|
Loading…
Reference in New Issue
Block a user