mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-11 14:33:52 +00:00
new: add roothash and state sync filter apis
This commit is contained in:
parent
236ec15f1b
commit
337d1bfb7c
17 changed files with 270 additions and 12 deletions
11
accounts/abi/bind/backends/bor_simulated.go
Normal file
11
accounts/abi/bind/backends/bor_simulated.go
Normal file
|
|
@ -0,0 +1,11 @@
|
||||||
|
package backends
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/ethereum/go-ethereum/core"
|
||||||
|
"github.com/ethereum/go-ethereum/event"
|
||||||
|
)
|
||||||
|
|
||||||
|
// SubscribeStateSyncEvent subscribes to state sync events
|
||||||
|
func (fb *filterBackend) SubscribeStateSyncEvent(ch chan<- core.StateSyncEvent) event.Subscription {
|
||||||
|
return fb.bc.SubscribeStateSyncEvent(ch)
|
||||||
|
}
|
||||||
|
|
@ -224,8 +224,7 @@ type Bor struct {
|
||||||
HeimdallClient IHeimdallClient
|
HeimdallClient IHeimdallClient
|
||||||
WithoutHeimdall bool
|
WithoutHeimdall bool
|
||||||
|
|
||||||
stateSyncFeed event.Feed
|
scope event.SubscriptionScope
|
||||||
scope event.SubscriptionScope
|
|
||||||
// The fields below are for testing only
|
// The fields below are for testing only
|
||||||
fakeDiff bool // Skip difficulty verifications
|
fakeDiff bool // Skip difficulty verifications
|
||||||
}
|
}
|
||||||
|
|
@ -653,6 +652,8 @@ func (c *Bor) Prepare(chain consensus.ChainHeaderReader, header *types.Header) e
|
||||||
// Finalize implements consensus.Engine, ensuring no uncles are set, nor block
|
// Finalize implements consensus.Engine, ensuring no uncles are set, nor block
|
||||||
// rewards given.
|
// rewards given.
|
||||||
func (c *Bor) Finalize(chain consensus.ChainHeaderReader, header *types.Header, state *state.StateDB, txs []*types.Transaction, uncles []*types.Header) {
|
func (c *Bor) Finalize(chain consensus.ChainHeaderReader, header *types.Header, state *state.StateDB, txs []*types.Transaction, uncles []*types.Header) {
|
||||||
|
stateSyncData := []*types.StateSyncData{}
|
||||||
|
|
||||||
var err error
|
var err error
|
||||||
headerNumber := header.Number.Uint64()
|
headerNumber := header.Number.Uint64()
|
||||||
if headerNumber%c.config.Sprint == 0 {
|
if headerNumber%c.config.Sprint == 0 {
|
||||||
|
|
@ -665,7 +666,7 @@ func (c *Bor) Finalize(chain consensus.ChainHeaderReader, header *types.Header,
|
||||||
|
|
||||||
if !c.WithoutHeimdall {
|
if !c.WithoutHeimdall {
|
||||||
// commit statees
|
// commit statees
|
||||||
_, err = c.CommitStates(state, header, cx)
|
stateSyncData, err = c.CommitStates(state, header, cx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Error while committing states", "error", err)
|
log.Error("Error while committing states", "error", err)
|
||||||
return
|
return
|
||||||
|
|
@ -676,11 +677,17 @@ func (c *Bor) Finalize(chain consensus.ChainHeaderReader, header *types.Header,
|
||||||
// No block rewards in PoA, so the state remains as is and uncles are dropped
|
// No block rewards in PoA, so the state remains as is and uncles are dropped
|
||||||
header.Root = state.IntermediateRoot(chain.Config().IsEIP158(header.Number))
|
header.Root = state.IntermediateRoot(chain.Config().IsEIP158(header.Number))
|
||||||
header.UncleHash = types.CalcUncleHash(nil)
|
header.UncleHash = types.CalcUncleHash(nil)
|
||||||
|
|
||||||
|
// Set state sync data to blockchain
|
||||||
|
bc := chain.(*core.BlockChain)
|
||||||
|
bc.SetStateSync(stateSyncData)
|
||||||
}
|
}
|
||||||
|
|
||||||
// FinalizeAndAssemble implements consensus.Engine, ensuring no uncles are set,
|
// FinalizeAndAssemble implements consensus.Engine, ensuring no uncles are set,
|
||||||
// nor block rewards given, and returns the final block.
|
// nor block rewards given, and returns the final block.
|
||||||
func (c *Bor) FinalizeAndAssemble(chain consensus.ChainHeaderReader, header *types.Header, state *state.StateDB, txs []*types.Transaction, uncles []*types.Header, receipts []*types.Receipt) (*types.Block, error) {
|
func (c *Bor) FinalizeAndAssemble(chain consensus.ChainHeaderReader, header *types.Header, state *state.StateDB, txs []*types.Transaction, uncles []*types.Header, receipts []*types.Receipt) (*types.Block, error) {
|
||||||
|
stateSyncData := []*types.StateSyncData{}
|
||||||
|
|
||||||
headerNumber := header.Number.Uint64()
|
headerNumber := header.Number.Uint64()
|
||||||
if headerNumber%c.config.Sprint == 0 {
|
if headerNumber%c.config.Sprint == 0 {
|
||||||
cx := chainContext{Chain: chain, Bor: c}
|
cx := chainContext{Chain: chain, Bor: c}
|
||||||
|
|
@ -694,7 +701,7 @@ func (c *Bor) FinalizeAndAssemble(chain consensus.ChainHeaderReader, header *typ
|
||||||
|
|
||||||
if !c.WithoutHeimdall {
|
if !c.WithoutHeimdall {
|
||||||
// commit states
|
// commit states
|
||||||
_, err = c.CommitStates(state, header, cx)
|
stateSyncData, err = c.CommitStates(state, header, cx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Error("Error while committing states", "error", err)
|
log.Error("Error while committing states", "error", err)
|
||||||
return nil, err
|
return nil, err
|
||||||
|
|
@ -708,6 +715,11 @@ func (c *Bor) FinalizeAndAssemble(chain consensus.ChainHeaderReader, header *typ
|
||||||
|
|
||||||
// Assemble block
|
// Assemble block
|
||||||
block := types.NewBlock(header, txs, nil, receipts, new(trie.Trie))
|
block := types.NewBlock(header, txs, nil, receipts, new(trie.Trie))
|
||||||
|
|
||||||
|
// set state sync
|
||||||
|
bc := chain.(*core.BlockChain)
|
||||||
|
bc.SetStateSync(stateSyncData)
|
||||||
|
|
||||||
// return the final block for sealing
|
// return the final block for sealing
|
||||||
return block, nil
|
return block, nil
|
||||||
}
|
}
|
||||||
|
|
@ -1090,8 +1102,8 @@ func (c *Bor) CommitStates(
|
||||||
state *state.StateDB,
|
state *state.StateDB,
|
||||||
header *types.Header,
|
header *types.Header,
|
||||||
chain chainContext,
|
chain chainContext,
|
||||||
) ([]*types.StateData, error) {
|
) ([]*types.StateSyncData, error) {
|
||||||
stateSyncs := make([]*types.StateData, 0)
|
stateSyncs := make([]*types.StateSyncData, 0)
|
||||||
number := header.Number.Uint64()
|
number := header.Number.Uint64()
|
||||||
_lastStateID, err := c.GenesisContractsClient.LastStateId(number - 1)
|
_lastStateID, err := c.GenesisContractsClient.LastStateId(number - 1)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
@ -1116,8 +1128,8 @@ func (c *Bor) CommitStates(
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|
||||||
stateData := types.StateData{
|
stateData := types.StateSyncData{
|
||||||
Did: eventRecord.ID,
|
ID: eventRecord.ID,
|
||||||
Contract: eventRecord.Contract,
|
Contract: eventRecord.Contract,
|
||||||
Data: hex.EncodeToString(eventRecord.Data),
|
Data: hex.EncodeToString(eventRecord.Data),
|
||||||
TxHash: eventRecord.TxHash,
|
TxHash: eventRecord.TxHash,
|
||||||
|
|
|
||||||
|
|
@ -211,6 +211,10 @@ type BlockChain struct {
|
||||||
shouldPreserve func(*types.Block) bool // Function used to determine whether should preserve the given block.
|
shouldPreserve func(*types.Block) bool // Function used to determine whether should preserve the given block.
|
||||||
terminateInsert func(common.Hash, uint64) bool // Testing hook used to terminate ancient receipt chain insertion.
|
terminateInsert func(common.Hash, uint64) bool // Testing hook used to terminate ancient receipt chain insertion.
|
||||||
writeLegacyJournal bool // Testing flag used to flush the snapshot journal in legacy format.
|
writeLegacyJournal bool // Testing flag used to flush the snapshot journal in legacy format.
|
||||||
|
|
||||||
|
// Bor related changes
|
||||||
|
stateSyncData []*types.StateSyncData // State sync data
|
||||||
|
stateSyncFeed event.Feed // State sync feed
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewBlockChain returns a fully initialised block chain using information
|
// NewBlockChain returns a fully initialised block chain using information
|
||||||
|
|
@ -1630,6 +1634,11 @@ func (bc *BlockChain) writeBlockWithState(block *types.Block, receipts []*types.
|
||||||
// event here.
|
// event here.
|
||||||
if emitHeadEvent {
|
if emitHeadEvent {
|
||||||
bc.chainHeadFeed.Send(ChainHeadEvent{Block: block})
|
bc.chainHeadFeed.Send(ChainHeadEvent{Block: block})
|
||||||
|
// BOR state sync feed related changes
|
||||||
|
for _, data := range bc.stateSyncData {
|
||||||
|
bc.stateSyncFeed.Send(StateSyncEvent{Data: data})
|
||||||
|
}
|
||||||
|
// BOR
|
||||||
}
|
}
|
||||||
} else {
|
} else {
|
||||||
bc.chainSideFeed.Send(ChainSideEvent{Block: block})
|
bc.chainSideFeed.Send(ChainSideEvent{Block: block})
|
||||||
|
|
@ -1887,6 +1896,12 @@ func (bc *BlockChain) insertChain(chain types.Blocks, verifySeals bool) (int, er
|
||||||
atomic.StoreUint32(&followupInterrupt, 1)
|
atomic.StoreUint32(&followupInterrupt, 1)
|
||||||
return it.index, err
|
return it.index, err
|
||||||
}
|
}
|
||||||
|
// BOR state sync feed related changes
|
||||||
|
for _, data := range bc.stateSyncData {
|
||||||
|
bc.stateSyncFeed.Send(StateSyncEvent{Data: data})
|
||||||
|
}
|
||||||
|
// BOR
|
||||||
|
|
||||||
// Update the metrics touched during block processing
|
// Update the metrics touched during block processing
|
||||||
accountReadTimer.Update(statedb.AccountReads) // Account reads are complete, we can mark them
|
accountReadTimer.Update(statedb.AccountReads) // Account reads are complete, we can mark them
|
||||||
storageReadTimer.Update(statedb.StorageReads) // Storage reads are complete, we can mark them
|
storageReadTimer.Update(statedb.StorageReads) // Storage reads are complete, we can mark them
|
||||||
|
|
@ -2555,3 +2570,17 @@ func (bc *BlockChain) SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscript
|
||||||
func (bc *BlockChain) SubscribeBlockProcessingEvent(ch chan<- bool) event.Subscription {
|
func (bc *BlockChain) SubscribeBlockProcessingEvent(ch chan<- bool) event.Subscription {
|
||||||
return bc.scope.Track(bc.blockProcFeed.Subscribe(ch))
|
return bc.scope.Track(bc.blockProcFeed.Subscribe(ch))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
//
|
||||||
|
// Bor related changes
|
||||||
|
//
|
||||||
|
|
||||||
|
// SetStateSync set sync data in state_data
|
||||||
|
func (bc *BlockChain) SetStateSync(stateData []*types.StateSyncData) {
|
||||||
|
bc.stateSyncData = stateData
|
||||||
|
}
|
||||||
|
|
||||||
|
// SubscribeStateSyncEvent registers a subscription of StateSyncEvent.
|
||||||
|
func (bc *BlockChain) SubscribeStateSyncEvent(ch chan<- StateSyncEvent) event.Subscription {
|
||||||
|
return bc.scope.Track(bc.stateSyncFeed.Subscribe(ch))
|
||||||
|
}
|
||||||
|
|
|
||||||
10
core/bor_events.go
Normal file
10
core/bor_events.go
Normal file
|
|
@ -0,0 +1,10 @@
|
||||||
|
package core
|
||||||
|
|
||||||
|
import (
|
||||||
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
|
)
|
||||||
|
|
||||||
|
// StateSyncEvent represents state sync events
|
||||||
|
type StateSyncEvent struct {
|
||||||
|
Data *types.StateSyncData
|
||||||
|
}
|
||||||
|
|
@ -2,9 +2,9 @@ package types
|
||||||
|
|
||||||
import "github.com/ethereum/go-ethereum/common"
|
import "github.com/ethereum/go-ethereum/common"
|
||||||
|
|
||||||
// StateData represents state received from Ethereum Blockchain
|
// StateSyncData represents state received from Ethereum Blockchain
|
||||||
type StateData struct {
|
type StateSyncData struct {
|
||||||
Did uint64
|
ID uint64
|
||||||
Contract common.Address
|
Contract common.Address
|
||||||
Data string
|
Data string
|
||||||
TxHash common.Hash
|
TxHash common.Hash
|
||||||
|
|
|
||||||
35
eth/bor_api_backend.go
Normal file
35
eth/bor_api_backend.go
Normal file
|
|
@ -0,0 +1,35 @@
|
||||||
|
package eth
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/consensus/bor"
|
||||||
|
"github.com/ethereum/go-ethereum/core"
|
||||||
|
"github.com/ethereum/go-ethereum/event"
|
||||||
|
)
|
||||||
|
|
||||||
|
// GetRootHash returns root hash for given start and end block
|
||||||
|
func (b *EthAPIBackend) GetRootHash(ctx context.Context, starBlockNr uint64, endBlockNr uint64) (string, error) {
|
||||||
|
var api *bor.API
|
||||||
|
for _, _api := range b.eth.Engine().APIs(b.eth.BlockChain()) {
|
||||||
|
if _api.Namespace == "bor" {
|
||||||
|
api = _api.Service.(*bor.API)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if api == nil {
|
||||||
|
return "", errors.New("Only available in Bor engine")
|
||||||
|
}
|
||||||
|
|
||||||
|
root, err := api.GetRootHash(starBlockNr, endBlockNr)
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
return root, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
// SubscribeStateSyncEvent subscribes to state sync event
|
||||||
|
func (b *EthAPIBackend) SubscribeStateSyncEvent(ch chan<- core.StateSyncEvent) event.Subscription {
|
||||||
|
return b.eth.BlockChain().SubscribeStateSyncEvent(ch)
|
||||||
|
}
|
||||||
43
eth/filters/bor_api.go
Normal file
43
eth/filters/bor_api.go
Normal file
|
|
@ -0,0 +1,43 @@
|
||||||
|
package filters
|
||||||
|
|
||||||
|
import (
|
||||||
|
"bytes"
|
||||||
|
"context"
|
||||||
|
|
||||||
|
ethereum "github.com/ethereum/go-ethereum"
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
|
"github.com/ethereum/go-ethereum/rpc"
|
||||||
|
)
|
||||||
|
|
||||||
|
// NewDeposits send a notification each time a new deposit received from bridge.
|
||||||
|
func (api *PublicFilterAPI) NewDeposits(ctx context.Context, crit ethereum.StateSyncFilter) (*rpc.Subscription, error) {
|
||||||
|
notifier, supported := rpc.NotifierFromContext(ctx)
|
||||||
|
if !supported {
|
||||||
|
return &rpc.Subscription{}, rpc.ErrNotificationsUnsupported
|
||||||
|
}
|
||||||
|
|
||||||
|
rpcSub := notifier.CreateSubscription()
|
||||||
|
go func() {
|
||||||
|
stateSyncData := make(chan *types.StateSyncData)
|
||||||
|
stateSyncSub := api.events.SubscribeNewDeposits(stateSyncData)
|
||||||
|
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case h := <-stateSyncData:
|
||||||
|
if crit.ID == h.ID || bytes.Compare(crit.Contract.Bytes(), h.Contract.Bytes()) == 0 ||
|
||||||
|
(crit.ID == 0 && crit.Contract == common.Address{}) {
|
||||||
|
notifier.Notify(rpcSub.ID, h)
|
||||||
|
}
|
||||||
|
case <-rpcSub.Err():
|
||||||
|
stateSyncSub.Unsubscribe()
|
||||||
|
return
|
||||||
|
case <-notifier.Closed():
|
||||||
|
stateSyncSub.Unsubscribe()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
return rpcSub, nil
|
||||||
|
}
|
||||||
32
eth/filters/bor_filter_system.go
Normal file
32
eth/filters/bor_filter_system.go
Normal file
|
|
@ -0,0 +1,32 @@
|
||||||
|
package filters
|
||||||
|
|
||||||
|
import (
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/core"
|
||||||
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
|
"github.com/ethereum/go-ethereum/rpc"
|
||||||
|
)
|
||||||
|
|
||||||
|
func (es *EventSystem) handleStateSyncEvent(filters filterIndex, ev core.StateSyncEvent) {
|
||||||
|
for _, f := range filters[StateSyncSubscription] {
|
||||||
|
f.stateSyncData <- ev.Data
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// SubscribeNewDeposits creates a subscription that writes details about the new state sync events (from mainchain to Bor)
|
||||||
|
func (es *EventSystem) SubscribeNewDeposits(data chan *types.StateSyncData) *Subscription {
|
||||||
|
sub := &subscription{
|
||||||
|
id: rpc.NewID(),
|
||||||
|
typ: StateSyncSubscription,
|
||||||
|
created: time.Now(),
|
||||||
|
logs: make(chan []*types.Log),
|
||||||
|
hashes: make(chan []common.Hash),
|
||||||
|
headers: make(chan *types.Header),
|
||||||
|
stateSyncData: data,
|
||||||
|
installed: make(chan struct{}),
|
||||||
|
err: make(chan error),
|
||||||
|
}
|
||||||
|
return es.subscribe(sub)
|
||||||
|
}
|
||||||
|
|
@ -45,6 +45,8 @@ type Backend interface {
|
||||||
|
|
||||||
BloomStatus() (uint64, uint64)
|
BloomStatus() (uint64, uint64)
|
||||||
ServiceFilter(ctx context.Context, session *bloombits.MatcherSession)
|
ServiceFilter(ctx context.Context, session *bloombits.MatcherSession)
|
||||||
|
|
||||||
|
SubscribeStateSyncEvent(ch chan<- core.StateSyncEvent) event.Subscription
|
||||||
}
|
}
|
||||||
|
|
||||||
// Filter can be used to retrieve and filter logs.
|
// Filter can be used to retrieve and filter logs.
|
||||||
|
|
|
||||||
|
|
@ -52,7 +52,9 @@ const (
|
||||||
PendingTransactionsSubscription
|
PendingTransactionsSubscription
|
||||||
// BlocksSubscription queries hashes for blocks that are imported
|
// BlocksSubscription queries hashes for blocks that are imported
|
||||||
BlocksSubscription
|
BlocksSubscription
|
||||||
// LastSubscription keeps track of the last index
|
// StateSyncSubscription to listen main chain state
|
||||||
|
StateSyncSubscription
|
||||||
|
// LastIndexSubscription keeps track of the last index
|
||||||
LastIndexSubscription
|
LastIndexSubscription
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -66,6 +68,8 @@ const (
|
||||||
logsChanSize = 10
|
logsChanSize = 10
|
||||||
// chainEvChanSize is the size of channel listening to ChainEvent.
|
// chainEvChanSize is the size of channel listening to ChainEvent.
|
||||||
chainEvChanSize = 10
|
chainEvChanSize = 10
|
||||||
|
// stateEvChanSize is the size of channel listening to StateSyncEvent.
|
||||||
|
stateEvChanSize = 10
|
||||||
)
|
)
|
||||||
|
|
||||||
type subscription struct {
|
type subscription struct {
|
||||||
|
|
@ -78,6 +82,8 @@ type subscription struct {
|
||||||
headers chan *types.Header
|
headers chan *types.Header
|
||||||
installed chan struct{} // closed when the filter is installed
|
installed chan struct{} // closed when the filter is installed
|
||||||
err chan error // closed when the filter is uninstalled
|
err chan error // closed when the filter is uninstalled
|
||||||
|
|
||||||
|
stateSyncData chan *types.StateSyncData
|
||||||
}
|
}
|
||||||
|
|
||||||
// EventSystem creates subscriptions, processes events and broadcasts them to the
|
// EventSystem creates subscriptions, processes events and broadcasts them to the
|
||||||
|
|
@ -102,6 +108,10 @@ type EventSystem struct {
|
||||||
pendingLogsCh chan []*types.Log // Channel to receive new log event
|
pendingLogsCh chan []*types.Log // Channel to receive new log event
|
||||||
rmLogsCh chan core.RemovedLogsEvent // Channel to receive removed log event
|
rmLogsCh chan core.RemovedLogsEvent // Channel to receive removed log event
|
||||||
chainCh chan core.ChainEvent // Channel to receive new chain event
|
chainCh chan core.ChainEvent // Channel to receive new chain event
|
||||||
|
|
||||||
|
// Bor related subscription and channels
|
||||||
|
stateSyncSub event.Subscription // Subscription for new state event
|
||||||
|
stateSyncCh chan core.StateSyncEvent // Channel to receive deposit state change event
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewEventSystem creates a new manager that listens for event on the given mux,
|
// NewEventSystem creates a new manager that listens for event on the given mux,
|
||||||
|
|
@ -121,6 +131,7 @@ func NewEventSystem(backend Backend, lightMode bool) *EventSystem {
|
||||||
rmLogsCh: make(chan core.RemovedLogsEvent, rmLogsChanSize),
|
rmLogsCh: make(chan core.RemovedLogsEvent, rmLogsChanSize),
|
||||||
pendingLogsCh: make(chan []*types.Log, logsChanSize),
|
pendingLogsCh: make(chan []*types.Log, logsChanSize),
|
||||||
chainCh: make(chan core.ChainEvent, chainEvChanSize),
|
chainCh: make(chan core.ChainEvent, chainEvChanSize),
|
||||||
|
stateSyncCh: make(chan core.StateSyncEvent, stateEvChanSize),
|
||||||
}
|
}
|
||||||
|
|
||||||
// Subscribe events
|
// Subscribe events
|
||||||
|
|
@ -129,6 +140,7 @@ func NewEventSystem(backend Backend, lightMode bool) *EventSystem {
|
||||||
m.rmLogsSub = m.backend.SubscribeRemovedLogsEvent(m.rmLogsCh)
|
m.rmLogsSub = m.backend.SubscribeRemovedLogsEvent(m.rmLogsCh)
|
||||||
m.chainSub = m.backend.SubscribeChainEvent(m.chainCh)
|
m.chainSub = m.backend.SubscribeChainEvent(m.chainCh)
|
||||||
m.pendingLogsSub = m.backend.SubscribePendingLogsEvent(m.pendingLogsCh)
|
m.pendingLogsSub = m.backend.SubscribePendingLogsEvent(m.pendingLogsCh)
|
||||||
|
m.stateSyncSub = m.backend.SubscribeStateSyncEvent(m.stateSyncCh)
|
||||||
|
|
||||||
// Make sure none of the subscriptions are empty
|
// Make sure none of the subscriptions are empty
|
||||||
if m.txsSub == nil || m.logsSub == nil || m.rmLogsSub == nil || m.chainSub == nil || m.pendingLogsSub == nil {
|
if m.txsSub == nil || m.logsSub == nil || m.rmLogsSub == nil || m.chainSub == nil || m.pendingLogsSub == nil {
|
||||||
|
|
@ -448,6 +460,7 @@ func (es *EventSystem) eventLoop() {
|
||||||
es.rmLogsSub.Unsubscribe()
|
es.rmLogsSub.Unsubscribe()
|
||||||
es.pendingLogsSub.Unsubscribe()
|
es.pendingLogsSub.Unsubscribe()
|
||||||
es.chainSub.Unsubscribe()
|
es.chainSub.Unsubscribe()
|
||||||
|
es.stateSyncSub.Unsubscribe()
|
||||||
}()
|
}()
|
||||||
|
|
||||||
index := make(filterIndex)
|
index := make(filterIndex)
|
||||||
|
|
@ -467,6 +480,8 @@ func (es *EventSystem) eventLoop() {
|
||||||
es.handlePendingLogs(index, ev)
|
es.handlePendingLogs(index, ev)
|
||||||
case ev := <-es.chainCh:
|
case ev := <-es.chainCh:
|
||||||
es.handleChainEvent(index, ev)
|
es.handleChainEvent(index, ev)
|
||||||
|
case ev := <-es.stateSyncCh:
|
||||||
|
es.handleStateSyncEvent(index, ev)
|
||||||
|
|
||||||
case f := <-es.install:
|
case f := <-es.install:
|
||||||
if f.typ == MinedAndPendingLogsSubscription {
|
if f.typ == MinedAndPendingLogsSubscription {
|
||||||
|
|
|
||||||
|
|
@ -47,6 +47,8 @@ type testBackend struct {
|
||||||
rmLogsFeed event.Feed
|
rmLogsFeed event.Feed
|
||||||
pendingLogsFeed event.Feed
|
pendingLogsFeed event.Feed
|
||||||
chainFeed event.Feed
|
chainFeed event.Feed
|
||||||
|
|
||||||
|
stateSyncFeed event.Feed
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *testBackend) ChainDb() ethdb.Database {
|
func (b *testBackend) ChainDb() ethdb.Database {
|
||||||
|
|
@ -152,6 +154,10 @@ func (b *testBackend) ServiceFilter(ctx context.Context, session *bloombits.Matc
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (b *testBackend) SubscribeStateSyncEvent(ch chan<- core.StateSyncEvent) event.Subscription {
|
||||||
|
return b.stateSyncFeed.Subscribe(ch)
|
||||||
|
}
|
||||||
|
|
||||||
// TestBlockSubscription tests if a block subscription returns block hashes for posted chain events.
|
// TestBlockSubscription tests if a block subscription returns block hashes for posted chain events.
|
||||||
// It creates multiple subscriptions:
|
// It creates multiple subscriptions:
|
||||||
// - one at the start and should receive all posted chain events and a second (blockHashes)
|
// - one at the start and should receive all posted chain events and a second (blockHashes)
|
||||||
|
|
|
||||||
14
ethclient/bor_ethclient.go
Normal file
14
ethclient/bor_ethclient.go
Normal file
|
|
@ -0,0 +1,14 @@
|
||||||
|
package ethclient
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
)
|
||||||
|
|
||||||
|
// GetRootHash returns the merkle root of the block headers
|
||||||
|
func (ec *Client) GetRootHash(ctx context.Context, startBlockNumber uint64, endBlockNumber uint64) (string, error) {
|
||||||
|
var rootHash string
|
||||||
|
if err := ec.c.CallContext(ctx, &rootHash, "eth_getRootHash", startBlockNumber, endBlockNumber); err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
return rootHash, nil
|
||||||
|
}
|
||||||
|
|
@ -209,3 +209,9 @@ type GasEstimator interface {
|
||||||
type PendingStateEventer interface {
|
type PendingStateEventer interface {
|
||||||
SubscribePendingTransactions(ctx context.Context, ch chan<- *types.Transaction) (Subscription, error)
|
SubscribePendingTransactions(ctx context.Context, ch chan<- *types.Transaction) (Subscription, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// StateSyncFilter state sync filter
|
||||||
|
type StateSyncFilter struct {
|
||||||
|
ID uint64
|
||||||
|
Contract common.Address
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -86,6 +86,10 @@ type Backend interface {
|
||||||
SubscribePendingLogsEvent(ch chan<- []*types.Log) event.Subscription
|
SubscribePendingLogsEvent(ch chan<- []*types.Log) event.Subscription
|
||||||
SubscribeRemovedLogsEvent(ch chan<- core.RemovedLogsEvent) event.Subscription
|
SubscribeRemovedLogsEvent(ch chan<- core.RemovedLogsEvent) event.Subscription
|
||||||
|
|
||||||
|
// Bor related APIs
|
||||||
|
SubscribeStateSyncEvent(ch chan<- core.StateSyncEvent) event.Subscription
|
||||||
|
GetRootHash(ctx context.Context, starBlockNr uint64, endBlockNr uint64) (string, error)
|
||||||
|
|
||||||
ChainConfig() *params.ChainConfig
|
ChainConfig() *params.ChainConfig
|
||||||
Engine() consensus.Engine
|
Engine() consensus.Engine
|
||||||
}
|
}
|
||||||
|
|
|
||||||
14
internal/ethapi/bor_api.go
Normal file
14
internal/ethapi/bor_api.go
Normal file
|
|
@ -0,0 +1,14 @@
|
||||||
|
package ethapi
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
)
|
||||||
|
|
||||||
|
// GetRootHash returns root hash for given start and end block
|
||||||
|
func (s *PublicBlockChainAPI) GetRootHash(ctx context.Context, starBlockNr uint64, endBlockNr uint64) (string, error) {
|
||||||
|
root, err := s.b.GetRootHash(ctx, starBlockNr, endBlockNr)
|
||||||
|
if err != nil {
|
||||||
|
return "", err
|
||||||
|
}
|
||||||
|
return root, nil
|
||||||
|
}
|
||||||
19
les/bor_api_backend.go
Normal file
19
les/bor_api_backend.go
Normal file
|
|
@ -0,0 +1,19 @@
|
||||||
|
package les
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"errors"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/core"
|
||||||
|
"github.com/ethereum/go-ethereum/event"
|
||||||
|
)
|
||||||
|
|
||||||
|
// GetRootHash returns root hash for given start and end block
|
||||||
|
func (b *LesApiBackend) GetRootHash(ctx context.Context, starBlockNr uint64, endBlockNr uint64) (string, error) {
|
||||||
|
return "", errors.New("Not implemented")
|
||||||
|
}
|
||||||
|
|
||||||
|
// SubscribeStateSyncEvent subscribe state sync event
|
||||||
|
func (b *LesApiBackend) SubscribeStateSyncEvent(ch chan<- core.StateSyncEvent) event.Subscription {
|
||||||
|
return b.eth.blockchain.SubscribeStateSyncEvent(ch)
|
||||||
|
}
|
||||||
|
|
@ -580,3 +580,9 @@ func (lc *LightChain) DisableCheckFreq() {
|
||||||
func (lc *LightChain) EnableCheckFreq() {
|
func (lc *LightChain) EnableCheckFreq() {
|
||||||
atomic.StoreInt32(&lc.disableCheckFreq, 0)
|
atomic.StoreInt32(&lc.disableCheckFreq, 0)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SubscribeStateSyncEvent implements the interface of filters.Backend
|
||||||
|
// LightChain does not send core.NewStateChangeSyncEvent, so return an empty subscription.
|
||||||
|
func (lc *LightChain) SubscribeStateSyncEvent(ch chan<- core.StateSyncEvent) event.Subscription {
|
||||||
|
return lc.scope.Track(new(event.Feed).Subscribe(ch))
|
||||||
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue