mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-07-23 13:16:42 +00:00
track delivered blocks
Signed-off-by: jsvisa <delweng@gmail.com>
This commit is contained in:
parent
99188a60b8
commit
b274824681
1 changed files with 59 additions and 13 deletions
|
|
@ -28,6 +28,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum"
|
"github.com/ethereum/go-ethereum"
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/common/hexutil"
|
"github.com/ethereum/go-ethereum/common/hexutil"
|
||||||
|
"github.com/ethereum/go-ethereum/common/lru"
|
||||||
"github.com/ethereum/go-ethereum/core/history"
|
"github.com/ethereum/go-ethereum/core/history"
|
||||||
"github.com/ethereum/go-ethereum/core/types"
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
"github.com/ethereum/go-ethereum/internal/ethapi"
|
"github.com/ethereum/go-ethereum/internal/ethapi"
|
||||||
|
|
@ -44,6 +45,11 @@ var (
|
||||||
errClientUnsubscribed = errors.New("client unsubscribed")
|
errClientUnsubscribed = errors.New("client unsubscribed")
|
||||||
)
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
// maxTrackedBlocks is the number of block hashes that will be tracked by subscription.
|
||||||
|
maxTrackedBlocks = 32 * 1024
|
||||||
|
)
|
||||||
|
|
||||||
// filter is a helper struct that holds meta information over the filter type
|
// filter is a helper struct that holds meta information over the filter type
|
||||||
// and associated subscription in the event system.
|
// and associated subscription in the event system.
|
||||||
type filter struct {
|
type filter struct {
|
||||||
|
|
@ -312,6 +318,7 @@ func (api *FilterAPI) liveLogs(notifier notifier, rpcSub *rpc.Subscription, crit
|
||||||
select {
|
select {
|
||||||
case logs := <-matchedLogs:
|
case logs := <-matchedLogs:
|
||||||
notifyLogsIf(notifier, rpcSub.ID, logs, nil)
|
notifyLogsIf(notifier, rpcSub.ID, logs, nil)
|
||||||
|
|
||||||
case <-rpcSub.Err(): // client send an unsubscribe request
|
case <-rpcSub.Err(): // client send an unsubscribe request
|
||||||
logsSub.Unsubscribe()
|
logsSub.Unsubscribe()
|
||||||
return
|
return
|
||||||
|
|
@ -360,14 +367,21 @@ func (api *FilterAPI) histLogs(notifier notifier, rpcSub *rpc.Subscription, from
|
||||||
var (
|
var (
|
||||||
// delivered is the block number of the last historical log delivered.
|
// delivered is the block number of the last historical log delivered.
|
||||||
delivered uint64
|
delivered uint64
|
||||||
|
|
||||||
// liveMode is true when either:
|
// liveMode is true when either:
|
||||||
// - all historical logs are delivered.
|
// - all historical logs are delivered.
|
||||||
// - or, during history processing a reorg is detected.
|
// - or, during history processing a reorg is detected.
|
||||||
liveMode bool
|
liveMode bool
|
||||||
|
|
||||||
// reorgedBlockHash is the block hash of the reorg point. It is set when
|
// reorgedBlockHash is the block hash of the reorg point. It is set when
|
||||||
// a reorg is detected in the future. It is used to detect if the history
|
// a reorg is detected in the future. It is used to detect if the history
|
||||||
// processor is sending stale logs.
|
// processor is sending stale logs.
|
||||||
reorgedBlockHash common.Hash
|
reorgedBlockHash common.Hash
|
||||||
|
|
||||||
|
// hashes is used to track the hashes of the blocks that have been delivered.
|
||||||
|
// It is used as a guard to prevent duplicate logs as well as inaccurate "removed"
|
||||||
|
// logs being delivered during a reorg.
|
||||||
|
hashes = lru.NewBasicLRU[common.Hash, struct{}](maxTrackedBlocks)
|
||||||
)
|
)
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
|
|
@ -394,12 +408,6 @@ func (api *FilterAPI) histLogs(notifier notifier, rpcSub *rpc.Subscription, from
|
||||||
// only send logs from the chain subscription.
|
// only send logs from the chain subscription.
|
||||||
logger.Info("Reorg detected", "reorgBlock", logs[0].BlockNumber, "delivered", delivered)
|
logger.Info("Reorg detected", "reorgBlock", logs[0].BlockNumber, "delivered", delivered)
|
||||||
liveMode = true
|
liveMode = true
|
||||||
// On reorg the chain will send first the removed logs.
|
|
||||||
// Send them up to the block we've delivered historical logs for.
|
|
||||||
notifyLogsIf(notifier, rpcSub.ID, logs, func(log *types.Log) bool {
|
|
||||||
return log.BlockNumber <= delivered
|
|
||||||
})
|
|
||||||
continue
|
|
||||||
}
|
}
|
||||||
if !liveMode {
|
if !liveMode {
|
||||||
if logs[0].Removed && reorgedBlockHash == (common.Hash{}) {
|
if logs[0].Removed && reorgedBlockHash == (common.Hash{}) {
|
||||||
|
|
@ -411,9 +419,8 @@ func (api *FilterAPI) histLogs(notifier notifier, rpcSub *rpc.Subscription, from
|
||||||
// - history is still being processed and blockchain sends logs from the tip.
|
// - history is still being processed and blockchain sends logs from the tip.
|
||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
// New logs are emitted for the whole new chain since the reorg block.
|
// Removed logs from reorged chain, replacing logs or logs from tip of the chain.
|
||||||
// Send them all.
|
notifyLogsIf(notifier, rpcSub.ID, logs, &hashes)
|
||||||
notifyLogsIf(notifier, rpcSub.ID, logs, nil)
|
|
||||||
|
|
||||||
case logs := <-histLogs:
|
case logs := <-histLogs:
|
||||||
if len(logs) == 0 {
|
if len(logs) == 0 {
|
||||||
|
|
@ -435,7 +442,7 @@ func (api *FilterAPI) histLogs(notifier notifier, rpcSub *rpc.Subscription, from
|
||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
notifyLogsIf(notifier, rpcSub.ID, logs, nil)
|
notifyLogsIf(notifier, rpcSub.ID, logs, &hashes)
|
||||||
// Assuming batch = all logs of a single block
|
// Assuming batch = all logs of a single block
|
||||||
delivered = logs[0].BlockNumber
|
delivered = logs[0].BlockNumber
|
||||||
|
|
||||||
|
|
@ -494,9 +501,48 @@ func (api *FilterAPI) doHistLogs(ctx context.Context, from int64, addrs []common
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func notifyLogsIf(notifier notifier, id rpc.ID, logs []*types.Log, cond func(log *types.Log) bool) {
|
// notifyLogsIf sends logs to the notifier if the condition is met.
|
||||||
for _, log := range logs {
|
// It assumes all logs of the same block are either all removed or all added.
|
||||||
if cond == nil || cond(log) {
|
func notifyLogsIf(notifier notifier, id rpc.ID, logs []*types.Log, hashes *lru.BasicLRU[common.Hash, struct{}]) {
|
||||||
|
// Iterate logs and batch them by block hash.
|
||||||
|
type batch struct {
|
||||||
|
start int
|
||||||
|
end int
|
||||||
|
hash common.Hash
|
||||||
|
removed bool
|
||||||
|
}
|
||||||
|
var (
|
||||||
|
batches = make([]batch, 0)
|
||||||
|
h common.Hash
|
||||||
|
)
|
||||||
|
for i, log := range logs {
|
||||||
|
if h == log.BlockHash {
|
||||||
|
// Skip logs of seen block
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if len(batches) > 0 {
|
||||||
|
batches[len(batches)-1].end = i
|
||||||
|
}
|
||||||
|
batches = append(batches, batch{start: i, hash: log.BlockHash, removed: log.Removed})
|
||||||
|
h = log.BlockHash
|
||||||
|
}
|
||||||
|
// Close off last batch.
|
||||||
|
if batches[len(batches)-1].end == 0 {
|
||||||
|
batches[len(batches)-1].end = len(logs)
|
||||||
|
}
|
||||||
|
for _, batch := range batches {
|
||||||
|
if hashes != nil {
|
||||||
|
// During reorgs it's possible that logs from the new chain have been delivered.
|
||||||
|
// Avoid sending removed logs from the old chain and duplicate logs from new chain.
|
||||||
|
if batch.removed && !hashes.Contains(batch.hash) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if !batch.removed && hashes.Contains(batch.hash) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
hashes.Add(batch.hash, struct{}{})
|
||||||
|
}
|
||||||
|
for _, log := range logs[batch.start:batch.end] {
|
||||||
log := log
|
log := log
|
||||||
notifier.Notify(id, &log)
|
notifier.Notify(id, &log)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue