mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-16 08:53:46 +00:00
Merge pull request #45 from maticnetwork/dev-subscribe-deposits
MAT-871: Subscribe to State Sync Events on Bor
This commit is contained in:
commit
86a8696a33
15 changed files with 161 additions and 29 deletions
|
|
@ -510,6 +510,9 @@ func (fb *filterBackend) SubscribeRemovedLogsEvent(ch chan<- core.RemovedLogsEve
|
||||||
func (fb *filterBackend) SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscription {
|
func (fb *filterBackend) SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscription {
|
||||||
return fb.bc.SubscribeLogsEvent(ch)
|
return fb.bc.SubscribeLogsEvent(ch)
|
||||||
}
|
}
|
||||||
|
func (fb *filterBackend) SubscribeStateEvent(ch chan<- core.NewStateChangeEvent) event.Subscription {
|
||||||
|
return fb.bc.SubscribeStateEvent(ch)
|
||||||
|
}
|
||||||
|
|
||||||
func (fb *filterBackend) BloomStatus() (uint64, uint64) { return 4096, 0 }
|
func (fb *filterBackend) BloomStatus() (uint64, uint64) { return 4096, 0 }
|
||||||
func (fb *filterBackend) ServiceFilter(ctx context.Context, ms *bloombits.MatcherSession) {
|
func (fb *filterBackend) ServiceFilter(ctx context.Context, ms *bloombits.MatcherSession) {
|
||||||
|
|
|
||||||
|
|
@ -30,6 +30,7 @@ import (
|
||||||
"github.com/maticnetwork/bor/core/vm"
|
"github.com/maticnetwork/bor/core/vm"
|
||||||
"github.com/maticnetwork/bor/crypto"
|
"github.com/maticnetwork/bor/crypto"
|
||||||
"github.com/maticnetwork/bor/ethdb"
|
"github.com/maticnetwork/bor/ethdb"
|
||||||
|
"github.com/maticnetwork/bor/event"
|
||||||
"github.com/maticnetwork/bor/internal/ethapi"
|
"github.com/maticnetwork/bor/internal/ethapi"
|
||||||
"github.com/maticnetwork/bor/log"
|
"github.com/maticnetwork/bor/log"
|
||||||
"github.com/maticnetwork/bor/params"
|
"github.com/maticnetwork/bor/params"
|
||||||
|
|
@ -245,6 +246,8 @@ type Bor struct {
|
||||||
stateReceiverABI abi.ABI
|
stateReceiverABI abi.ABI
|
||||||
HeimdallClient IHeimdallClient
|
HeimdallClient IHeimdallClient
|
||||||
|
|
||||||
|
stateDataFeed event.Feed
|
||||||
|
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
|
||||||
}
|
}
|
||||||
|
|
@ -1201,6 +1204,16 @@ func (c *Bor) CommitStates(
|
||||||
"txHash", eventRecord.TxHash,
|
"txHash", eventRecord.TxHash,
|
||||||
"chainID", eventRecord.ChainID,
|
"chainID", eventRecord.ChainID,
|
||||||
)
|
)
|
||||||
|
stateData := types.StateData{
|
||||||
|
Did: eventRecord.ID,
|
||||||
|
Contract: eventRecord.Contract,
|
||||||
|
Data: hex.EncodeToString(eventRecord.Data),
|
||||||
|
TxHash: eventRecord.TxHash,
|
||||||
|
}
|
||||||
|
|
||||||
|
go func() {
|
||||||
|
c.stateDataFeed.Send(core.NewStateChangeEvent{StateData: &stateData})
|
||||||
|
}()
|
||||||
|
|
||||||
recordBytes, err := rlp.EncodeToBytes(eventRecord)
|
recordBytes, err := rlp.EncodeToBytes(eventRecord)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|
@ -1226,6 +1239,11 @@ func (c *Bor) CommitStates(
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SubscribeStateEvent registers a subscription of ChainSideEvent.
|
||||||
|
func (c *Bor) SubscribeStateEvent(ch chan<- core.NewStateChangeEvent) event.Subscription {
|
||||||
|
return c.scope.Track(c.stateDataFeed.Subscribe(ch))
|
||||||
|
}
|
||||||
|
|
||||||
func (c *Bor) SetHeimdallClient(h IHeimdallClient) {
|
func (c *Bor) SetHeimdallClient(h IHeimdallClient) {
|
||||||
c.HeimdallClient = h
|
c.HeimdallClient = h
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -27,6 +27,7 @@ import (
|
||||||
"sync/atomic"
|
"sync/atomic"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
|
"github.com/hashicorp/golang-lru"
|
||||||
"github.com/maticnetwork/bor/common"
|
"github.com/maticnetwork/bor/common"
|
||||||
"github.com/maticnetwork/bor/common/mclock"
|
"github.com/maticnetwork/bor/common/mclock"
|
||||||
"github.com/maticnetwork/bor/common/prque"
|
"github.com/maticnetwork/bor/common/prque"
|
||||||
|
|
@ -42,7 +43,6 @@ import (
|
||||||
"github.com/maticnetwork/bor/params"
|
"github.com/maticnetwork/bor/params"
|
||||||
"github.com/maticnetwork/bor/rlp"
|
"github.com/maticnetwork/bor/rlp"
|
||||||
"github.com/maticnetwork/bor/trie"
|
"github.com/maticnetwork/bor/trie"
|
||||||
"github.com/hashicorp/golang-lru"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
|
|
@ -142,6 +142,7 @@ type BlockChain struct {
|
||||||
chainHeadFeed event.Feed
|
chainHeadFeed event.Feed
|
||||||
logsFeed event.Feed
|
logsFeed event.Feed
|
||||||
blockProcFeed event.Feed
|
blockProcFeed event.Feed
|
||||||
|
stateDataFeed event.Feed
|
||||||
scope event.SubscriptionScope
|
scope event.SubscriptionScope
|
||||||
genesisBlock *types.Block
|
genesisBlock *types.Block
|
||||||
|
|
||||||
|
|
@ -2177,6 +2178,11 @@ func (bc *BlockChain) SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscript
|
||||||
return bc.scope.Track(bc.logsFeed.Subscribe(ch))
|
return bc.scope.Track(bc.logsFeed.Subscribe(ch))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SubscribeStateEvent registers a subscription of ChainSideEvent.
|
||||||
|
func (bc *BlockChain) SubscribeStateEvent(ch chan<- NewStateChangeEvent) event.Subscription {
|
||||||
|
return bc.scope.Track(bc.stateDataFeed.Subscribe(ch))
|
||||||
|
}
|
||||||
|
|
||||||
// SubscribeBlockProcessingEvent registers a subscription of bool where true means
|
// SubscribeBlockProcessingEvent registers a subscription of bool where true means
|
||||||
// block processing has started while false means it has stopped.
|
// block processing has started while false means it has stopped.
|
||||||
func (bc *BlockChain) SubscribeBlockProcessingEvent(ch chan<- bool) event.Subscription {
|
func (bc *BlockChain) SubscribeBlockProcessingEvent(ch chan<- bool) event.Subscription {
|
||||||
|
|
|
||||||
|
|
@ -41,6 +41,10 @@ type ChainEvent struct {
|
||||||
Logs []*types.Log
|
Logs []*types.Log
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type NewStateChangeEvent struct {
|
||||||
|
StateData *types.StateData
|
||||||
|
}
|
||||||
|
|
||||||
type ChainSideEvent struct {
|
type ChainSideEvent struct {
|
||||||
Block *types.Block
|
Block *types.Block
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -60,6 +60,14 @@ type txdata struct {
|
||||||
Hash *common.Hash `json:"hash" rlp:"-"`
|
Hash *common.Hash `json:"hash" rlp:"-"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// State represents state received from Ethereum Blockchain
|
||||||
|
type StateData struct {
|
||||||
|
Did uint64
|
||||||
|
Contract common.Address
|
||||||
|
Data string
|
||||||
|
TxHash common.Hash
|
||||||
|
}
|
||||||
|
|
||||||
type txdataMarshaling struct {
|
type txdataMarshaling struct {
|
||||||
AccountNonce hexutil.Uint64
|
AccountNonce hexutil.Uint64
|
||||||
Price *hexutil.Big
|
Price *hexutil.Big
|
||||||
|
|
|
||||||
|
|
@ -24,6 +24,7 @@ import (
|
||||||
"github.com/maticnetwork/bor/accounts"
|
"github.com/maticnetwork/bor/accounts"
|
||||||
"github.com/maticnetwork/bor/common"
|
"github.com/maticnetwork/bor/common"
|
||||||
"github.com/maticnetwork/bor/common/math"
|
"github.com/maticnetwork/bor/common/math"
|
||||||
|
"github.com/maticnetwork/bor/consensus/bor"
|
||||||
"github.com/maticnetwork/bor/core"
|
"github.com/maticnetwork/bor/core"
|
||||||
"github.com/maticnetwork/bor/core/bloombits"
|
"github.com/maticnetwork/bor/core/bloombits"
|
||||||
"github.com/maticnetwork/bor/core/rawdb"
|
"github.com/maticnetwork/bor/core/rawdb"
|
||||||
|
|
@ -155,6 +156,11 @@ func (b *EthAPIBackend) SubscribeChainSideEvent(ch chan<- core.ChainSideEvent) e
|
||||||
return b.eth.BlockChain().SubscribeChainSideEvent(ch)
|
return b.eth.BlockChain().SubscribeChainSideEvent(ch)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (b *EthAPIBackend) SubscribeStateEvent(ch chan<- core.NewStateChangeEvent) event.Subscription {
|
||||||
|
engine := b.eth.Engine()
|
||||||
|
return engine.(*bor.Bor).SubscribeStateEvent(ch)
|
||||||
|
}
|
||||||
|
|
||||||
func (b *EthAPIBackend) SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscription {
|
func (b *EthAPIBackend) SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscription {
|
||||||
return b.eth.BlockChain().SubscribeLogsEvent(ch)
|
return b.eth.BlockChain().SubscribeLogsEvent(ch)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -17,6 +17,7 @@
|
||||||
package filters
|
package filters
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
|
|
@ -215,7 +216,6 @@ func (api *PublicFilterAPI) NewHeads(ctx context.Context) (*rpc.Subscription, er
|
||||||
go func() {
|
go func() {
|
||||||
headers := make(chan *types.Header)
|
headers := make(chan *types.Header)
|
||||||
headersSub := api.events.SubscribeNewHeads(headers)
|
headersSub := api.events.SubscribeNewHeads(headers)
|
||||||
|
|
||||||
for {
|
for {
|
||||||
select {
|
select {
|
||||||
case h := <-headers:
|
case h := <-headers:
|
||||||
|
|
@ -233,6 +233,38 @@ func (api *PublicFilterAPI) NewHeads(ctx context.Context) (*rpc.Subscription, er
|
||||||
return rpcSub, nil
|
return rpcSub, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// NewDeposits send a notification each time a new deposit received from bridge.
|
||||||
|
func (api *PublicFilterAPI) NewDeposits(ctx context.Context, crit ethereum.FilterState) (*rpc.Subscription, error) {
|
||||||
|
notifier, supported := rpc.NotifierFromContext(ctx)
|
||||||
|
if !supported {
|
||||||
|
return &rpc.Subscription{}, rpc.ErrNotificationsUnsupported
|
||||||
|
}
|
||||||
|
|
||||||
|
rpcSub := notifier.CreateSubscription()
|
||||||
|
go func() {
|
||||||
|
stateData := make(chan *types.StateData)
|
||||||
|
stateDataSub := api.events.SubscribeNewDeposits(stateData)
|
||||||
|
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case h := <-stateData:
|
||||||
|
if crit.Did == h.Did || bytes.Compare(crit.Contract.Bytes(), h.Contract.Bytes()) == 0 ||
|
||||||
|
(crit.Did == 0 && crit.Contract == common.Address{}) {
|
||||||
|
notifier.Notify(rpcSub.ID, h)
|
||||||
|
}
|
||||||
|
case <-rpcSub.Err():
|
||||||
|
stateDataSub.Unsubscribe()
|
||||||
|
return
|
||||||
|
case <-notifier.Closed():
|
||||||
|
stateDataSub.Unsubscribe()
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
return rpcSub, nil
|
||||||
|
}
|
||||||
|
|
||||||
// Logs creates a subscription that fires for all new log that match the given filter criteria.
|
// Logs creates a subscription that fires for all new log that match the given filter criteria.
|
||||||
func (api *PublicFilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subscription, error) {
|
func (api *PublicFilterAPI) Logs(ctx context.Context, crit FilterCriteria) (*rpc.Subscription, error) {
|
||||||
notifier, supported := rpc.NotifierFromContext(ctx)
|
notifier, supported := rpc.NotifierFromContext(ctx)
|
||||||
|
|
|
||||||
|
|
@ -40,6 +40,7 @@ type Backend interface {
|
||||||
|
|
||||||
SubscribeNewTxsEvent(chan<- core.NewTxsEvent) event.Subscription
|
SubscribeNewTxsEvent(chan<- core.NewTxsEvent) event.Subscription
|
||||||
SubscribeChainEvent(ch chan<- core.ChainEvent) event.Subscription
|
SubscribeChainEvent(ch chan<- core.ChainEvent) event.Subscription
|
||||||
|
SubscribeStateEvent(ch chan<- core.NewStateChangeEvent) event.Subscription
|
||||||
SubscribeRemovedLogsEvent(ch chan<- core.RemovedLogsEvent) event.Subscription
|
SubscribeRemovedLogsEvent(ch chan<- core.RemovedLogsEvent) event.Subscription
|
||||||
SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscription
|
SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscription
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -53,6 +53,10 @@ const (
|
||||||
PendingTransactionsSubscription
|
PendingTransactionsSubscription
|
||||||
// BlocksSubscription queries hashes for blocks that are imported
|
// BlocksSubscription queries hashes for blocks that are imported
|
||||||
BlocksSubscription
|
BlocksSubscription
|
||||||
|
|
||||||
|
//StateSubscription to listen main chain state
|
||||||
|
StateSubscription
|
||||||
|
|
||||||
// LastSubscription keeps track of the last index
|
// LastSubscription keeps track of the last index
|
||||||
LastIndexSubscription
|
LastIndexSubscription
|
||||||
)
|
)
|
||||||
|
|
@ -68,6 +72,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 ChainEvent.
|
||||||
|
stateEvChanSize = 10
|
||||||
)
|
)
|
||||||
|
|
||||||
var (
|
var (
|
||||||
|
|
@ -82,6 +88,7 @@ type subscription struct {
|
||||||
logs chan []*types.Log
|
logs chan []*types.Log
|
||||||
hashes chan []common.Hash
|
hashes chan []common.Hash
|
||||||
headers chan *types.Header
|
headers chan *types.Header
|
||||||
|
stateData chan *types.StateData
|
||||||
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
|
||||||
}
|
}
|
||||||
|
|
@ -99,15 +106,18 @@ type EventSystem struct {
|
||||||
logsSub event.Subscription // Subscription for new log event
|
logsSub event.Subscription // Subscription for new log event
|
||||||
rmLogsSub event.Subscription // Subscription for removed log event
|
rmLogsSub event.Subscription // Subscription for removed log event
|
||||||
chainSub event.Subscription // Subscription for new chain event
|
chainSub event.Subscription // Subscription for new chain event
|
||||||
|
stateSub event.Subscription // Subscription for new state change event
|
||||||
pendingLogSub *event.TypeMuxSubscription // Subscription for pending log event
|
pendingLogSub *event.TypeMuxSubscription // Subscription for pending log event
|
||||||
|
|
||||||
// Channels
|
// Channels
|
||||||
install chan *subscription // install filter for event notification
|
install chan *subscription // install filter for event notification
|
||||||
uninstall chan *subscription // remove filter for event notification
|
uninstall chan *subscription // remove filter for event notification
|
||||||
txsCh chan core.NewTxsEvent // Channel to receive new transactions event
|
txsCh chan core.NewTxsEvent // Channel to receive new transactions event
|
||||||
logsCh chan []*types.Log // Channel to receive new log event
|
logsCh 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
|
||||||
|
stateCh chan core.NewStateChangeEvent // 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,
|
||||||
|
|
@ -127,6 +137,7 @@ func NewEventSystem(mux *event.TypeMux, backend Backend, lightMode bool) *EventS
|
||||||
logsCh: make(chan []*types.Log, logsChanSize),
|
logsCh: make(chan []*types.Log, logsChanSize),
|
||||||
rmLogsCh: make(chan core.RemovedLogsEvent, rmLogsChanSize),
|
rmLogsCh: make(chan core.RemovedLogsEvent, rmLogsChanSize),
|
||||||
chainCh: make(chan core.ChainEvent, chainEvChanSize),
|
chainCh: make(chan core.ChainEvent, chainEvChanSize),
|
||||||
|
stateCh: make(chan core.NewStateChangeEvent, stateEvChanSize),
|
||||||
}
|
}
|
||||||
|
|
||||||
// Subscribe events
|
// Subscribe events
|
||||||
|
|
@ -134,12 +145,13 @@ func NewEventSystem(mux *event.TypeMux, backend Backend, lightMode bool) *EventS
|
||||||
m.logsSub = m.backend.SubscribeLogsEvent(m.logsCh)
|
m.logsSub = m.backend.SubscribeLogsEvent(m.logsCh)
|
||||||
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.stateSub = m.backend.SubscribeStateEvent(m.stateCh)
|
||||||
// TODO(rjl493456442): use feed to subscribe pending log event
|
// TODO(rjl493456442): use feed to subscribe pending log event
|
||||||
m.pendingLogSub = m.mux.Subscribe(core.PendingLogsEvent{})
|
m.pendingLogSub = m.mux.Subscribe(core.PendingLogsEvent{})
|
||||||
|
|
||||||
// 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 ||
|
if m.txsSub == nil || m.logsSub == nil || m.rmLogsSub == nil || m.chainSub == nil ||
|
||||||
m.pendingLogSub.Closed() {
|
m.stateSub == nil || m.pendingLogSub.Closed() {
|
||||||
log.Crit("Subscribe for event system failed")
|
log.Crit("Subscribe for event system failed")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -292,6 +304,23 @@ func (es *EventSystem) SubscribeNewHeads(headers chan *types.Header) *Subscripti
|
||||||
logs: make(chan []*types.Log),
|
logs: make(chan []*types.Log),
|
||||||
hashes: make(chan []common.Hash),
|
hashes: make(chan []common.Hash),
|
||||||
headers: headers,
|
headers: headers,
|
||||||
|
stateData: make(chan *types.StateData),
|
||||||
|
installed: make(chan struct{}),
|
||||||
|
err: make(chan error),
|
||||||
|
}
|
||||||
|
return es.subscribe(sub)
|
||||||
|
}
|
||||||
|
|
||||||
|
// SubscribeNewDeposits creates a subscription that writes details about the new state sync events (from mainchain to Bor)
|
||||||
|
func (es *EventSystem) SubscribeNewDeposits(stateData chan *types.StateData) *Subscription {
|
||||||
|
sub := &subscription{
|
||||||
|
id: rpc.NewID(),
|
||||||
|
typ: StateSubscription,
|
||||||
|
created: time.Now(),
|
||||||
|
logs: make(chan []*types.Log),
|
||||||
|
hashes: make(chan []common.Hash),
|
||||||
|
headers: make(chan *types.Header),
|
||||||
|
stateData: stateData,
|
||||||
installed: make(chan struct{}),
|
installed: make(chan struct{}),
|
||||||
err: make(chan error),
|
err: make(chan error),
|
||||||
}
|
}
|
||||||
|
|
@ -355,6 +384,10 @@ func (es *EventSystem) broadcast(filters filterIndex, ev interface{}) {
|
||||||
for _, f := range filters[PendingTransactionsSubscription] {
|
for _, f := range filters[PendingTransactionsSubscription] {
|
||||||
f.hashes <- hashes
|
f.hashes <- hashes
|
||||||
}
|
}
|
||||||
|
case core.NewStateChangeEvent:
|
||||||
|
for _, f := range filters[StateSubscription] {
|
||||||
|
f.stateData <- e.StateData
|
||||||
|
}
|
||||||
case core.ChainEvent:
|
case core.ChainEvent:
|
||||||
for _, f := range filters[BlocksSubscription] {
|
for _, f := range filters[BlocksSubscription] {
|
||||||
f.headers <- e.Block.Header()
|
f.headers <- e.Block.Header()
|
||||||
|
|
@ -471,6 +504,8 @@ func (es *EventSystem) eventLoop() {
|
||||||
es.broadcast(index, ev)
|
es.broadcast(index, ev)
|
||||||
case ev := <-es.chainCh:
|
case ev := <-es.chainCh:
|
||||||
es.broadcast(index, ev)
|
es.broadcast(index, ev)
|
||||||
|
case ev := <-es.stateCh:
|
||||||
|
es.broadcast(index, ev)
|
||||||
case ev, active := <-es.pendingLogSub.Chan():
|
case ev, active := <-es.pendingLogSub.Chan():
|
||||||
if !active { // system stopped
|
if !active { // system stopped
|
||||||
return
|
return
|
||||||
|
|
|
||||||
|
|
@ -24,7 +24,7 @@ import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"math/big"
|
"math/big"
|
||||||
|
|
||||||
"github.com/maticnetwork/bor"
|
ethereum "github.com/maticnetwork/bor"
|
||||||
"github.com/maticnetwork/bor/common"
|
"github.com/maticnetwork/bor/common"
|
||||||
"github.com/maticnetwork/bor/common/hexutil"
|
"github.com/maticnetwork/bor/common/hexutil"
|
||||||
"github.com/maticnetwork/bor/core/types"
|
"github.com/maticnetwork/bor/core/types"
|
||||||
|
|
@ -324,6 +324,11 @@ func (ec *Client) SubscribeNewHead(ctx context.Context, ch chan<- *types.Header)
|
||||||
return ec.c.EthSubscribe(ctx, ch, "newHeads")
|
return ec.c.EthSubscribe(ctx, ch, "newHeads")
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SubscribeNewDeposit subscribes to new state sync events
|
||||||
|
func (ec *Client) SubscribeNewDeposit(ctx context.Context, ch chan<- *types.StateData) (ethereum.Subscription, error) {
|
||||||
|
return ec.c.EthSubscribe(ctx, ch, "newDeposits", nil)
|
||||||
|
}
|
||||||
|
|
||||||
// State Access
|
// State Access
|
||||||
|
|
||||||
// NetworkID returns the network ID (also known as the chain ID) for this chain.
|
// NetworkID returns the network ID (also known as the chain ID) for this chain.
|
||||||
|
|
|
||||||
|
|
@ -25,7 +25,7 @@ import (
|
||||||
"testing"
|
"testing"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
"github.com/maticnetwork/bor"
|
bor "github.com/maticnetwork/bor"
|
||||||
"github.com/maticnetwork/bor/common"
|
"github.com/maticnetwork/bor/common"
|
||||||
"github.com/maticnetwork/bor/consensus/ethash"
|
"github.com/maticnetwork/bor/consensus/ethash"
|
||||||
"github.com/maticnetwork/bor/core"
|
"github.com/maticnetwork/bor/core"
|
||||||
|
|
@ -39,17 +39,17 @@ import (
|
||||||
|
|
||||||
// Verify that Client implements the ethereum interfaces.
|
// Verify that Client implements the ethereum interfaces.
|
||||||
var (
|
var (
|
||||||
_ = ethereum.ChainReader(&Client{})
|
_ = bor.ChainReader(&Client{})
|
||||||
_ = ethereum.TransactionReader(&Client{})
|
_ = bor.TransactionReader(&Client{})
|
||||||
_ = ethereum.ChainStateReader(&Client{})
|
_ = bor.ChainStateReader(&Client{})
|
||||||
_ = ethereum.ChainSyncReader(&Client{})
|
_ = bor.ChainSyncReader(&Client{})
|
||||||
_ = ethereum.ContractCaller(&Client{})
|
_ = bor.ContractCaller(&Client{})
|
||||||
_ = ethereum.GasEstimator(&Client{})
|
_ = bor.GasEstimator(&Client{})
|
||||||
_ = ethereum.GasPricer(&Client{})
|
_ = bor.GasPricer(&Client{})
|
||||||
_ = ethereum.LogFilterer(&Client{})
|
_ = bor.LogFilterer(&Client{})
|
||||||
_ = ethereum.PendingStateReader(&Client{})
|
_ = bor.PendingStateReader(&Client{})
|
||||||
// _ = ethereum.PendingStateEventer(&Client{})
|
// _ = bor.PendingStateEventer(&Client{})
|
||||||
_ = ethereum.PendingContractCaller(&Client{})
|
_ = bor.PendingContractCaller(&Client{})
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestToFilterArg(t *testing.T) {
|
func TestToFilterArg(t *testing.T) {
|
||||||
|
|
@ -63,13 +63,13 @@ func TestToFilterArg(t *testing.T) {
|
||||||
|
|
||||||
for _, testCase := range []struct {
|
for _, testCase := range []struct {
|
||||||
name string
|
name string
|
||||||
input ethereum.FilterQuery
|
input bor.FilterQuery
|
||||||
output interface{}
|
output interface{}
|
||||||
err error
|
err error
|
||||||
}{
|
}{
|
||||||
{
|
{
|
||||||
"without BlockHash",
|
"without BlockHash",
|
||||||
ethereum.FilterQuery{
|
bor.FilterQuery{
|
||||||
Addresses: addresses,
|
Addresses: addresses,
|
||||||
FromBlock: big.NewInt(1),
|
FromBlock: big.NewInt(1),
|
||||||
ToBlock: big.NewInt(2),
|
ToBlock: big.NewInt(2),
|
||||||
|
|
@ -85,7 +85,7 @@ func TestToFilterArg(t *testing.T) {
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"with nil fromBlock and nil toBlock",
|
"with nil fromBlock and nil toBlock",
|
||||||
ethereum.FilterQuery{
|
bor.FilterQuery{
|
||||||
Addresses: addresses,
|
Addresses: addresses,
|
||||||
Topics: [][]common.Hash{},
|
Topics: [][]common.Hash{},
|
||||||
},
|
},
|
||||||
|
|
@ -99,7 +99,7 @@ func TestToFilterArg(t *testing.T) {
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"with blockhash",
|
"with blockhash",
|
||||||
ethereum.FilterQuery{
|
bor.FilterQuery{
|
||||||
Addresses: addresses,
|
Addresses: addresses,
|
||||||
BlockHash: &blockHash,
|
BlockHash: &blockHash,
|
||||||
Topics: [][]common.Hash{},
|
Topics: [][]common.Hash{},
|
||||||
|
|
@ -113,7 +113,7 @@ func TestToFilterArg(t *testing.T) {
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"with blockhash and from block",
|
"with blockhash and from block",
|
||||||
ethereum.FilterQuery{
|
bor.FilterQuery{
|
||||||
Addresses: addresses,
|
Addresses: addresses,
|
||||||
BlockHash: &blockHash,
|
BlockHash: &blockHash,
|
||||||
FromBlock: big.NewInt(1),
|
FromBlock: big.NewInt(1),
|
||||||
|
|
@ -124,7 +124,7 @@ func TestToFilterArg(t *testing.T) {
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"with blockhash and to block",
|
"with blockhash and to block",
|
||||||
ethereum.FilterQuery{
|
bor.FilterQuery{
|
||||||
Addresses: addresses,
|
Addresses: addresses,
|
||||||
BlockHash: &blockHash,
|
BlockHash: &blockHash,
|
||||||
ToBlock: big.NewInt(1),
|
ToBlock: big.NewInt(1),
|
||||||
|
|
@ -135,7 +135,7 @@ func TestToFilterArg(t *testing.T) {
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"with blockhash and both from / to block",
|
"with blockhash and both from / to block",
|
||||||
ethereum.FilterQuery{
|
bor.FilterQuery{
|
||||||
Addresses: addresses,
|
Addresses: addresses,
|
||||||
BlockHash: &blockHash,
|
BlockHash: &blockHash,
|
||||||
FromBlock: big.NewInt(1),
|
FromBlock: big.NewInt(1),
|
||||||
|
|
|
||||||
|
|
@ -150,6 +150,11 @@ type FilterQuery struct {
|
||||||
Topics [][]common.Hash
|
Topics [][]common.Hash
|
||||||
}
|
}
|
||||||
|
|
||||||
|
type FilterState struct {
|
||||||
|
Did uint64
|
||||||
|
Contract common.Address
|
||||||
|
}
|
||||||
|
|
||||||
// LogFilterer provides access to contract log events using a one-off query or continuous
|
// LogFilterer provides access to contract log events using a one-off query or continuous
|
||||||
// event subscription.
|
// event subscription.
|
||||||
//
|
//
|
||||||
|
|
|
||||||
|
|
@ -59,6 +59,7 @@ type Backend interface {
|
||||||
GetTd(blockHash common.Hash) *big.Int
|
GetTd(blockHash common.Hash) *big.Int
|
||||||
GetEVM(ctx context.Context, msg core.Message, state *state.StateDB, header *types.Header) (*vm.EVM, func() error, error)
|
GetEVM(ctx context.Context, msg core.Message, state *state.StateDB, header *types.Header) (*vm.EVM, func() error, error)
|
||||||
SubscribeChainEvent(ch chan<- core.ChainEvent) event.Subscription
|
SubscribeChainEvent(ch chan<- core.ChainEvent) event.Subscription
|
||||||
|
SubscribeStateEvent(ch chan<- core.NewStateChangeEvent) event.Subscription
|
||||||
SubscribeChainHeadEvent(ch chan<- core.ChainHeadEvent) event.Subscription
|
SubscribeChainHeadEvent(ch chan<- core.ChainHeadEvent) event.Subscription
|
||||||
SubscribeChainSideEvent(ch chan<- core.ChainSideEvent) event.Subscription
|
SubscribeChainSideEvent(ch chan<- core.ChainSideEvent) event.Subscription
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -24,6 +24,7 @@ import (
|
||||||
"github.com/maticnetwork/bor/accounts"
|
"github.com/maticnetwork/bor/accounts"
|
||||||
"github.com/maticnetwork/bor/common"
|
"github.com/maticnetwork/bor/common"
|
||||||
"github.com/maticnetwork/bor/common/math"
|
"github.com/maticnetwork/bor/common/math"
|
||||||
|
"github.com/maticnetwork/bor/consensus/bor"
|
||||||
"github.com/maticnetwork/bor/core"
|
"github.com/maticnetwork/bor/core"
|
||||||
"github.com/maticnetwork/bor/core/bloombits"
|
"github.com/maticnetwork/bor/core/bloombits"
|
||||||
"github.com/maticnetwork/bor/core/rawdb"
|
"github.com/maticnetwork/bor/core/rawdb"
|
||||||
|
|
@ -164,6 +165,11 @@ func (b *LesApiBackend) SubscribeChainSideEvent(ch chan<- core.ChainSideEvent) e
|
||||||
return b.eth.blockchain.SubscribeChainSideEvent(ch)
|
return b.eth.blockchain.SubscribeChainSideEvent(ch)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (b *LesApiBackend) SubscribeStateEvent(ch chan<- core.NewStateChangeEvent) event.Subscription {
|
||||||
|
engine := b.eth.Engine()
|
||||||
|
return engine.(*bor.Bor).SubscribeStateEvent(ch)
|
||||||
|
}
|
||||||
|
|
||||||
func (b *LesApiBackend) SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscription {
|
func (b *LesApiBackend) SubscribeLogsEvent(ch chan<- []*types.Log) event.Subscription {
|
||||||
return b.eth.blockchain.SubscribeLogsEvent(ch)
|
return b.eth.blockchain.SubscribeLogsEvent(ch)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -28,6 +28,7 @@ import (
|
||||||
"github.com/maticnetwork/bor/common/hexutil"
|
"github.com/maticnetwork/bor/common/hexutil"
|
||||||
"github.com/maticnetwork/bor/common/mclock"
|
"github.com/maticnetwork/bor/common/mclock"
|
||||||
"github.com/maticnetwork/bor/consensus"
|
"github.com/maticnetwork/bor/consensus"
|
||||||
|
"github.com/maticnetwork/bor/consensus/bor"
|
||||||
"github.com/maticnetwork/bor/core"
|
"github.com/maticnetwork/bor/core"
|
||||||
"github.com/maticnetwork/bor/core/bloombits"
|
"github.com/maticnetwork/bor/core/bloombits"
|
||||||
"github.com/maticnetwork/bor/core/rawdb"
|
"github.com/maticnetwork/bor/core/rawdb"
|
||||||
|
|
@ -59,6 +60,7 @@ type LightEthereum struct {
|
||||||
peers *peerSet
|
peers *peerSet
|
||||||
txPool *light.TxPool
|
txPool *light.TxPool
|
||||||
blockchain *light.LightChain
|
blockchain *light.LightChain
|
||||||
|
bor *bor.Bor
|
||||||
serverPool *serverPool
|
serverPool *serverPool
|
||||||
reqDist *requestDistributor
|
reqDist *requestDistributor
|
||||||
retriever *retrieveManager
|
retriever *retrieveManager
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue