mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-17 01:13:45 +00:00
light, les: CHT and bloom trie post processors
This commit is contained in:
parent
976133a71a
commit
b844a630b0
13 changed files with 496 additions and 185 deletions
|
|
@ -54,6 +54,7 @@ type LesServer interface {
|
||||||
Start(srvr *p2p.Server)
|
Start(srvr *p2p.Server)
|
||||||
Stop()
|
Stop()
|
||||||
Protocols() []p2p.Protocol
|
Protocols() []p2p.Protocol
|
||||||
|
SetBloomBitsIndexer(bbIndexer *core.ChainIndexer)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Ethereum implements the Ethereum full node service.
|
// Ethereum implements the Ethereum full node service.
|
||||||
|
|
@ -95,6 +96,7 @@ type Ethereum struct {
|
||||||
|
|
||||||
func (s *Ethereum) AddLesServer(ls LesServer) {
|
func (s *Ethereum) AddLesServer(ls LesServer) {
|
||||||
s.lesServer = ls
|
s.lesServer = ls
|
||||||
|
ls.SetBloomBitsIndexer(s.bloomIndexer)
|
||||||
}
|
}
|
||||||
|
|
||||||
// New creates a new Ethereum object (including the
|
// New creates a new Ethereum object (including the
|
||||||
|
|
@ -154,7 +156,7 @@ func New(ctx *node.ServiceContext, config *Config) (*Ethereum, error) {
|
||||||
eth.blockchain.SetHead(compat.RewindTo)
|
eth.blockchain.SetHead(compat.RewindTo)
|
||||||
core.WriteChainConfig(chainDb, genesisHash, chainConfig)
|
core.WriteChainConfig(chainDb, genesisHash, chainConfig)
|
||||||
}
|
}
|
||||||
eth.bloomIndexer.Start(eth.blockchain.CurrentHeader(), eth.blockchain.SubscribeChainEvent)
|
eth.bloomIndexer.Start(eth.blockchain)
|
||||||
|
|
||||||
if config.TxPool.Journal != "" {
|
if config.TxPool.Journal != "" {
|
||||||
config.TxPool.Journal = ctx.ResolvePath(config.TxPool.Journal)
|
config.TxPool.Journal = ctx.ResolvePath(config.TxPool.Journal)
|
||||||
|
|
|
||||||
|
|
@ -174,8 +174,12 @@ func (b *LesApiBackend) AccountManager() *accounts.Manager {
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *LesApiBackend) BloomStatus() (uint64, uint64) {
|
func (b *LesApiBackend) BloomStatus() (uint64, uint64) {
|
||||||
return params.BloomBitsBlocks, 0
|
sections, _, _ := b.eth.bbIndexer.Sections()
|
||||||
|
return light.BloomTrieFrequency, sections
|
||||||
}
|
}
|
||||||
|
|
||||||
func (b *LesApiBackend) ServiceFilter(ctx context.Context, session *bloombits.MatcherSession) {
|
func (b *LesApiBackend) ServiceFilter(ctx context.Context, session *bloombits.MatcherSession) {
|
||||||
|
for i := 0; i < bloomFilterThreads; i++ {
|
||||||
|
go session.Multiplex(bloomRetrievalBatch, bloomRetrievalWait, b.eth.bloomRequests)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -61,6 +61,8 @@ type LightEthereum struct {
|
||||||
// DB interfaces
|
// DB interfaces
|
||||||
chainDb ethdb.Database // Block chain database
|
chainDb ethdb.Database // Block chain database
|
||||||
|
|
||||||
|
bbIndexer, chtIndexer, bltIndexer *core.ChainIndexer
|
||||||
|
|
||||||
ApiBackend *LesApiBackend
|
ApiBackend *LesApiBackend
|
||||||
|
|
||||||
eventMux *event.TypeMux
|
eventMux *event.TypeMux
|
||||||
|
|
@ -87,7 +89,7 @@ func New(ctx *node.ServiceContext, config *eth.Config) (*LightEthereum, error) {
|
||||||
peers := newPeerSet()
|
peers := newPeerSet()
|
||||||
quitSync := make(chan struct{})
|
quitSync := make(chan struct{})
|
||||||
|
|
||||||
eth := &LightEthereum{
|
leth := &LightEthereum{
|
||||||
chainConfig: chainConfig,
|
chainConfig: chainConfig,
|
||||||
chainDb: chainDb,
|
chainDb: chainDb,
|
||||||
eventMux: ctx.EventMux,
|
eventMux: ctx.EventMux,
|
||||||
|
|
@ -97,33 +99,38 @@ func New(ctx *node.ServiceContext, config *eth.Config) (*LightEthereum, error) {
|
||||||
engine: eth.CreateConsensusEngine(ctx, config, chainConfig, chainDb),
|
engine: eth.CreateConsensusEngine(ctx, config, chainConfig, chainDb),
|
||||||
shutdownChan: make(chan bool),
|
shutdownChan: make(chan bool),
|
||||||
networkId: config.NetworkId,
|
networkId: config.NetworkId,
|
||||||
|
bbIndexer: eth.NewBloomBitsProcessor(chainDb, light.BloomTrieFrequency),
|
||||||
|
chtIndexer: light.NewChtIndexer(chainDb, true),
|
||||||
|
bltIndexer: light.NewBloomTrieIndexer(chainDb, true),
|
||||||
}
|
}
|
||||||
|
|
||||||
eth.relay = NewLesTxRelay(peers, eth.reqDist)
|
leth.relay = NewLesTxRelay(peers, leth.reqDist)
|
||||||
eth.serverPool = newServerPool(chainDb, quitSync, ð.wg)
|
leth.serverPool = newServerPool(chainDb, quitSync, &leth.wg)
|
||||||
eth.retriever = newRetrieveManager(peers, eth.reqDist, eth.serverPool)
|
leth.retriever = newRetrieveManager(peers, leth.reqDist, leth.serverPool)
|
||||||
eth.odr = NewLesOdr(chainDb, eth.retriever)
|
leth.odr = NewLesOdr(chainDb, leth.chtIndexer, leth.bltIndexer, leth.bbIndexer, leth.retriever)
|
||||||
if eth.blockchain, err = light.NewLightChain(eth.odr, eth.chainConfig, eth.engine); err != nil {
|
if leth.blockchain, err = light.NewLightChain(leth.odr, leth.chainConfig, leth.engine); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
// Rewind the chain in case of an incompatible config upgrade.
|
// Rewind the chain in case of an incompatible config upgrade.
|
||||||
if compat, ok := genesisErr.(*params.ConfigCompatError); ok {
|
if compat, ok := genesisErr.(*params.ConfigCompatError); ok {
|
||||||
log.Warn("Rewinding chain to upgrade configuration", "err", compat)
|
log.Warn("Rewinding chain to upgrade configuration", "err", compat)
|
||||||
eth.blockchain.SetHead(compat.RewindTo)
|
leth.blockchain.SetHead(compat.RewindTo)
|
||||||
core.WriteChainConfig(chainDb, genesisHash, chainConfig)
|
core.WriteChainConfig(chainDb, genesisHash, chainConfig)
|
||||||
}
|
}
|
||||||
|
|
||||||
eth.txPool = light.NewTxPool(eth.chainConfig, eth.blockchain, eth.relay)
|
//leth.bbIndexer.Start(leth.blockchain)
|
||||||
if eth.protocolManager, err = NewProtocolManager(eth.chainConfig, true, config.NetworkId, eth.eventMux, eth.engine, eth.peers, eth.blockchain, nil, chainDb, eth.odr, eth.relay, quitSync, ð.wg); err != nil {
|
|
||||||
|
leth.txPool = light.NewTxPool(leth.chainConfig, leth.blockchain, leth.relay)
|
||||||
|
if leth.protocolManager, err = NewProtocolManager(leth.chainConfig, true, config.NetworkId, leth.eventMux, leth.engine, leth.peers, leth.blockchain, nil, chainDb, leth.odr, leth.relay, quitSync, &leth.wg); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
eth.ApiBackend = &LesApiBackend{eth, nil}
|
leth.ApiBackend = &LesApiBackend{leth, nil}
|
||||||
gpoParams := config.GPO
|
gpoParams := config.GPO
|
||||||
if gpoParams.Default == nil {
|
if gpoParams.Default == nil {
|
||||||
gpoParams.Default = config.GasPrice
|
gpoParams.Default = config.GasPrice
|
||||||
}
|
}
|
||||||
eth.ApiBackend.gpo = gasprice.NewOracle(eth.ApiBackend, gpoParams)
|
leth.ApiBackend.gpo = gasprice.NewOracle(leth.ApiBackend, gpoParams)
|
||||||
return eth, nil
|
return leth, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func lesTopic(genesisHash common.Hash) discv5.Topic {
|
func lesTopic(genesisHash common.Hash) discv5.Topic {
|
||||||
|
|
@ -211,6 +218,9 @@ func (s *LightEthereum) Start(srvr *p2p.Server) error {
|
||||||
// Ethereum protocol.
|
// Ethereum protocol.
|
||||||
func (s *LightEthereum) Stop() error {
|
func (s *LightEthereum) Stop() error {
|
||||||
s.odr.Stop()
|
s.odr.Stop()
|
||||||
|
s.bbIndexer.Close()
|
||||||
|
s.chtIndexer.Close()
|
||||||
|
s.bltIndexer.Close()
|
||||||
s.blockchain.Stop()
|
s.blockchain.Stop()
|
||||||
s.protocolManager.Stop()
|
s.protocolManager.Stop()
|
||||||
s.txPool.Stop()
|
s.txPool.Stop()
|
||||||
|
|
|
||||||
|
|
@ -825,14 +825,16 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
|
||||||
if reject(uint64(reqCnt), MaxPPTProofsFetch) {
|
if reject(uint64(reqCnt), MaxPPTProofsFetch) {
|
||||||
return errResp(ErrRequestRejected, "")
|
return errResp(ErrRequestRejected, "")
|
||||||
}
|
}
|
||||||
|
trieDb := ethdb.NewTable(pm.chainDb, light.ChtTablePrefix)
|
||||||
for _, req := range req.Reqs {
|
for _, req := range req.Reqs {
|
||||||
if bytes >= softResponseLimit {
|
if bytes >= softResponseLimit {
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
|
|
||||||
if header := pm.blockchain.GetHeaderByNumber(req.BlockNum); header != nil {
|
if header := pm.blockchain.GetHeaderByNumber(req.BlockNum); header != nil {
|
||||||
if root := getChtRoot(pm.chainDb, req.ChtNum); root != (common.Hash{}) {
|
sectionHead := core.GetCanonicalHash(pm.chainDb, (req.ChtNum+1)*light.ChtV1Frequency-1)
|
||||||
if tr, _ := trie.New(root, pm.chainDb); tr != nil {
|
if root := light.GetChtRoot(pm.chainDb, req.ChtNum, sectionHead); root != (common.Hash{}) {
|
||||||
|
if tr, _ := trie.New(root, trieDb); tr != nil {
|
||||||
var encNumber [8]byte
|
var encNumber [8]byte
|
||||||
binary.BigEndian.PutUint64(encNumber[:], req.BlockNum)
|
binary.BigEndian.PutUint64(encNumber[:], req.BlockNum)
|
||||||
var proof light.NodeList
|
var proof light.NodeList
|
||||||
|
|
@ -997,9 +999,11 @@ func (pm *ProtocolManager) handleMsg(p *peer) error {
|
||||||
func (pm *ProtocolManager) getPPT(id uint, idx uint64) (common.Hash, string) {
|
func (pm *ProtocolManager) getPPT(id uint, idx uint64) (common.Hash, string) {
|
||||||
switch id {
|
switch id {
|
||||||
case PPTChain:
|
case PPTChain:
|
||||||
return light.GetChtRoot(pm.chainDb, idx), light.ChtTablePrefix
|
sectionHead := core.GetCanonicalHash(pm.chainDb, (idx+1)*light.ChtFrequency-1)
|
||||||
|
return light.GetChtV2Root(pm.chainDb, idx, sectionHead), light.ChtTablePrefix
|
||||||
case PPTBloomBits:
|
case PPTBloomBits:
|
||||||
return light.GetBloomTrieRoot(pm.chainDb, idx), light.BloomTrieTablePrefix
|
sectionHead := core.GetCanonicalHash(pm.chainDb, (idx+1)*light.BloomTrieFrequency-1)
|
||||||
|
return light.GetBloomTrieRoot(pm.chainDb, idx, sectionHead), light.BloomTrieTablePrefix
|
||||||
}
|
}
|
||||||
return common.Hash{}, ""
|
return common.Hash{}, ""
|
||||||
}
|
}
|
||||||
|
|
|
||||||
36
les/odr.go
36
les/odr.go
|
|
@ -19,6 +19,7 @@ package les
|
||||||
import (
|
import (
|
||||||
"context"
|
"context"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/core"
|
||||||
"github.com/ethereum/go-ethereum/ethdb"
|
"github.com/ethereum/go-ethereum/ethdb"
|
||||||
"github.com/ethereum/go-ethereum/light"
|
"github.com/ethereum/go-ethereum/light"
|
||||||
"github.com/ethereum/go-ethereum/log"
|
"github.com/ethereum/go-ethereum/log"
|
||||||
|
|
@ -26,27 +27,48 @@ import (
|
||||||
|
|
||||||
// LesOdr implements light.OdrBackend
|
// LesOdr implements light.OdrBackend
|
||||||
type LesOdr struct {
|
type LesOdr struct {
|
||||||
db ethdb.Database
|
db ethdb.Database
|
||||||
stop chan struct{}
|
chtIndexer, bltIndexer, bloomIndexer *core.ChainIndexer
|
||||||
retriever *retrieveManager
|
retriever *retrieveManager
|
||||||
|
stop chan struct{}
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewLesOdr(db ethdb.Database, retriever *retrieveManager) *LesOdr {
|
func NewLesOdr(db ethdb.Database, chtIndexer, bltIndexer, bloomIndexer *core.ChainIndexer, retriever *retrieveManager) *LesOdr {
|
||||||
return &LesOdr{
|
return &LesOdr{
|
||||||
db: db,
|
db: db,
|
||||||
retriever: retriever,
|
chtIndexer: chtIndexer,
|
||||||
stop: make(chan struct{}),
|
bltIndexer: bltIndexer,
|
||||||
|
bloomIndexer: bloomIndexer,
|
||||||
|
retriever: retriever,
|
||||||
|
stop: make(chan struct{}),
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Stop cancels all pending retrievals
|
||||||
func (odr *LesOdr) Stop() {
|
func (odr *LesOdr) Stop() {
|
||||||
close(odr.stop)
|
close(odr.stop)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Database returns the backing database
|
||||||
func (odr *LesOdr) Database() ethdb.Database {
|
func (odr *LesOdr) Database() ethdb.Database {
|
||||||
return odr.db
|
return odr.db
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ChtIndexer returns the CHT chain indexer
|
||||||
|
func (odr *LesOdr) ChtIndexer() *core.ChainIndexer {
|
||||||
|
return odr.chtIndexer
|
||||||
|
}
|
||||||
|
|
||||||
|
// BltIndexer returns the bloom trie chain indexer
|
||||||
|
func (odr *LesOdr) BltIndexer() *core.ChainIndexer {
|
||||||
|
return odr.bltIndexer
|
||||||
|
}
|
||||||
|
|
||||||
|
// BloomIndexer returns the bloombits chain indexer
|
||||||
|
func (odr *LesOdr) BloomIndexer() *core.ChainIndexer {
|
||||||
|
return odr.bloomIndexer
|
||||||
|
}
|
||||||
|
|
||||||
const (
|
const (
|
||||||
MsgBlockBodies = iota
|
MsgBlockBodies = iota
|
||||||
MsgCode
|
MsgCode
|
||||||
|
|
|
||||||
|
|
@ -362,7 +362,7 @@ func (r *ChtRequest) CanSend(peer *peer) bool {
|
||||||
peer.lock.RLock()
|
peer.lock.RLock()
|
||||||
defer peer.lock.RUnlock()
|
defer peer.lock.RUnlock()
|
||||||
|
|
||||||
return peer.headInfo.Number >= light.ChtConfirmations && r.ChtNum <= (peer.headInfo.Number-light.ChtConfirmations)/light.ChtFrequency
|
return peer.headInfo.Number >= light.PPTConfirmations && r.ChtNum <= (peer.headInfo.Number-light.PPTConfirmations)/light.ChtFrequency
|
||||||
}
|
}
|
||||||
|
|
||||||
// Request sends an ODR request to the LES network (implementation of LesOdrRequest)
|
// Request sends an ODR request to the LES network (implementation of LesOdrRequest)
|
||||||
|
|
@ -481,7 +481,7 @@ func (r *BloomRequest) CanSend(peer *peer) bool {
|
||||||
if peer.version < lpv2 {
|
if peer.version < lpv2 {
|
||||||
return false
|
return false
|
||||||
}
|
}
|
||||||
return peer.headInfo.Number >= light.BloomTrieConfirmations && r.BltNum <= (peer.headInfo.Number-light.BloomTrieConfirmations)/light.BloomTrieFrequency
|
return peer.headInfo.Number >= light.PPTConfirmations && r.BltNum <= (peer.headInfo.Number-light.PPTConfirmations)/light.BloomTrieFrequency
|
||||||
}
|
}
|
||||||
|
|
||||||
// Request sends an ODR request to the LES network (implementation of LesOdrRequest)
|
// Request sends an ODR request to the LES network (implementation of LesOdrRequest)
|
||||||
|
|
|
||||||
|
|
@ -29,6 +29,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/core/state"
|
"github.com/ethereum/go-ethereum/core/state"
|
||||||
"github.com/ethereum/go-ethereum/core/types"
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
"github.com/ethereum/go-ethereum/core/vm"
|
"github.com/ethereum/go-ethereum/core/vm"
|
||||||
|
"github.com/ethereum/go-ethereum/eth"
|
||||||
"github.com/ethereum/go-ethereum/ethdb"
|
"github.com/ethereum/go-ethereum/ethdb"
|
||||||
"github.com/ethereum/go-ethereum/light"
|
"github.com/ethereum/go-ethereum/light"
|
||||||
"github.com/ethereum/go-ethereum/params"
|
"github.com/ethereum/go-ethereum/params"
|
||||||
|
|
@ -154,7 +155,7 @@ func testOdr(t *testing.T, protocol int, expFail uint64, fn odrTestFn) {
|
||||||
rm := newRetrieveManager(peers, dist, nil)
|
rm := newRetrieveManager(peers, dist, nil)
|
||||||
db, _ := ethdb.NewMemDatabase()
|
db, _ := ethdb.NewMemDatabase()
|
||||||
ldb, _ := ethdb.NewMemDatabase()
|
ldb, _ := ethdb.NewMemDatabase()
|
||||||
odr := NewLesOdr(ldb, rm)
|
odr := NewLesOdr(ldb, light.NewChtIndexer(db, true), light.NewBloomTrieIndexer(db, true), eth.NewBloomIndexer(db, light.BloomTrieFrequency), rm)
|
||||||
pm := newTestProtocolManagerMust(t, false, 4, testChainGen, nil, nil, db)
|
pm := newTestProtocolManagerMust(t, false, 4, testChainGen, nil, nil, db)
|
||||||
lpm := newTestProtocolManagerMust(t, true, 0, nil, peers, odr, ldb)
|
lpm := newTestProtocolManagerMust(t, true, 0, nil, peers, odr, ldb)
|
||||||
_, err1, lpeer, err2 := newTestPeerPair("peer", protocol, pm, lpm)
|
_, err1, lpeer, err2 := newTestPeerPair("peer", protocol, pm, lpm)
|
||||||
|
|
|
||||||
|
|
@ -24,6 +24,7 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/core"
|
"github.com/ethereum/go-ethereum/core"
|
||||||
"github.com/ethereum/go-ethereum/crypto"
|
"github.com/ethereum/go-ethereum/crypto"
|
||||||
|
"github.com/ethereum/go-ethereum/eth"
|
||||||
"github.com/ethereum/go-ethereum/ethdb"
|
"github.com/ethereum/go-ethereum/ethdb"
|
||||||
"github.com/ethereum/go-ethereum/light"
|
"github.com/ethereum/go-ethereum/light"
|
||||||
)
|
)
|
||||||
|
|
@ -73,7 +74,7 @@ func testAccess(t *testing.T, protocol int, fn accessTestFn) {
|
||||||
rm := newRetrieveManager(peers, dist, nil)
|
rm := newRetrieveManager(peers, dist, nil)
|
||||||
db, _ := ethdb.NewMemDatabase()
|
db, _ := ethdb.NewMemDatabase()
|
||||||
ldb, _ := ethdb.NewMemDatabase()
|
ldb, _ := ethdb.NewMemDatabase()
|
||||||
odr := NewLesOdr(ldb, rm)
|
odr := NewLesOdr(ldb, light.NewChtIndexer(db, true), light.NewBloomTrieIndexer(db, true), eth.NewBloomIndexer(db, light.BloomTrieFrequency), rm)
|
||||||
|
|
||||||
pm := newTestProtocolManagerMust(t, false, 4, testChainGen, nil, nil, db)
|
pm := newTestProtocolManagerMust(t, false, 4, testChainGen, nil, nil, db)
|
||||||
lpm := newTestProtocolManagerMust(t, true, 0, nil, peers, odr, ldb)
|
lpm := newTestProtocolManagerMust(t, true, 0, nil, peers, odr, ldb)
|
||||||
|
|
|
||||||
110
les/server.go
110
les/server.go
|
|
@ -21,7 +21,6 @@ import (
|
||||||
"encoding/binary"
|
"encoding/binary"
|
||||||
"math"
|
"math"
|
||||||
"sync"
|
"sync"
|
||||||
"time"
|
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/core"
|
"github.com/ethereum/go-ethereum/core"
|
||||||
|
|
@ -34,7 +33,6 @@ import (
|
||||||
"github.com/ethereum/go-ethereum/p2p"
|
"github.com/ethereum/go-ethereum/p2p"
|
||||||
"github.com/ethereum/go-ethereum/p2p/discv5"
|
"github.com/ethereum/go-ethereum/p2p/discv5"
|
||||||
"github.com/ethereum/go-ethereum/rlp"
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
"github.com/ethereum/go-ethereum/trie"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
type LesServer struct {
|
type LesServer struct {
|
||||||
|
|
@ -44,6 +42,8 @@ type LesServer struct {
|
||||||
defParams *flowcontrol.ServerParams
|
defParams *flowcontrol.ServerParams
|
||||||
lesTopic discv5.Topic
|
lesTopic discv5.Topic
|
||||||
quitSync chan struct{}
|
quitSync chan struct{}
|
||||||
|
|
||||||
|
chtIndexer, bltIndexer *core.ChainIndexer
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewLesServer(eth *eth.Ethereum, config *eth.Config) (*LesServer, error) {
|
func NewLesServer(eth *eth.Ethereum, config *eth.Config) (*LesServer, error) {
|
||||||
|
|
@ -58,7 +58,10 @@ func NewLesServer(eth *eth.Ethereum, config *eth.Config) (*LesServer, error) {
|
||||||
protocolManager: pm,
|
protocolManager: pm,
|
||||||
quitSync: quitSync,
|
quitSync: quitSync,
|
||||||
lesTopic: lesTopic(eth.BlockChain().Genesis().Hash()),
|
lesTopic: lesTopic(eth.BlockChain().Genesis().Hash()),
|
||||||
|
chtIndexer: light.NewChtIndexer(eth.ChainDb(), false),
|
||||||
|
bltIndexer: light.NewBloomTrieIndexer(eth.ChainDb(), false),
|
||||||
}
|
}
|
||||||
|
srv.chtIndexer.Start(eth.BlockChain())
|
||||||
pm.server = srv
|
pm.server = srv
|
||||||
|
|
||||||
srv.defParams = &flowcontrol.ServerParams{
|
srv.defParams = &flowcontrol.ServerParams{
|
||||||
|
|
@ -86,8 +89,14 @@ func (s *LesServer) Start(srvr *p2p.Server) {
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *LesServer) SetBloomBitsIndexer(bbIndexer *core.ChainIndexer) {
|
||||||
|
bbIndexer.AddChildIndexer(s.bltIndexer)
|
||||||
|
}
|
||||||
|
|
||||||
// Stop stops the LES service
|
// Stop stops the LES service
|
||||||
func (s *LesServer) Stop() {
|
func (s *LesServer) Stop() {
|
||||||
|
s.chtIndexer.Close()
|
||||||
|
// bloom trie indexer is closed by parent bloombits indexer
|
||||||
s.fcCostStats.store()
|
s.fcCostStats.store()
|
||||||
s.fcManager.Stop()
|
s.fcManager.Stop()
|
||||||
go func() {
|
go func() {
|
||||||
|
|
@ -273,10 +282,7 @@ func (pm *ProtocolManager) blockLoop() {
|
||||||
pm.wg.Add(1)
|
pm.wg.Add(1)
|
||||||
headCh := make(chan core.ChainHeadEvent, 10)
|
headCh := make(chan core.ChainHeadEvent, 10)
|
||||||
headSub := pm.blockchain.SubscribeChainHeadEvent(headCh)
|
headSub := pm.blockchain.SubscribeChainHeadEvent(headCh)
|
||||||
newCht := make(chan struct{}, 10)
|
|
||||||
newCht <- struct{}{}
|
|
||||||
go func() {
|
go func() {
|
||||||
var mu sync.Mutex
|
|
||||||
var lastHead *types.Header
|
var lastHead *types.Header
|
||||||
lastBroadcastTd := common.Big0
|
lastBroadcastTd := common.Big0
|
||||||
for {
|
for {
|
||||||
|
|
@ -308,17 +314,6 @@ func (pm *ProtocolManager) blockLoop() {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
newCht <- struct{}{}
|
|
||||||
case <-newCht:
|
|
||||||
go func() {
|
|
||||||
mu.Lock()
|
|
||||||
more := makeCht(pm.chainDb)
|
|
||||||
mu.Unlock()
|
|
||||||
if more {
|
|
||||||
time.Sleep(time.Millisecond * 10)
|
|
||||||
newCht <- struct{}{}
|
|
||||||
}
|
|
||||||
}()
|
|
||||||
case <-pm.quitSync:
|
case <-pm.quitSync:
|
||||||
headSub.Unsubscribe()
|
headSub.Unsubscribe()
|
||||||
pm.wg.Done()
|
pm.wg.Done()
|
||||||
|
|
@ -327,86 +322,3 @@ func (pm *ProtocolManager) blockLoop() {
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
}
|
}
|
||||||
|
|
||||||
var (
|
|
||||||
lastChtKey = []byte("LastChtNumber") // chtNum (uint64 big endian)
|
|
||||||
chtPrefix = []byte("cht") // chtPrefix + chtNum (uint64 big endian) -> trie root hash
|
|
||||||
)
|
|
||||||
|
|
||||||
func getChtRoot(db ethdb.Database, num uint64) common.Hash {
|
|
||||||
var encNumber [8]byte
|
|
||||||
binary.BigEndian.PutUint64(encNumber[:], num)
|
|
||||||
data, _ := db.Get(append(chtPrefix, encNumber[:]...))
|
|
||||||
return common.BytesToHash(data)
|
|
||||||
}
|
|
||||||
|
|
||||||
func storeChtRoot(db ethdb.Database, num uint64, root common.Hash) {
|
|
||||||
var encNumber [8]byte
|
|
||||||
binary.BigEndian.PutUint64(encNumber[:], num)
|
|
||||||
db.Put(append(chtPrefix, encNumber[:]...), root[:])
|
|
||||||
}
|
|
||||||
|
|
||||||
func makeCht(db ethdb.Database) bool {
|
|
||||||
headHash := core.GetHeadBlockHash(db)
|
|
||||||
headNum := core.GetBlockNumber(db, headHash)
|
|
||||||
|
|
||||||
var newChtNum uint64
|
|
||||||
if headNum > light.ChtConfirmations {
|
|
||||||
newChtNum = (headNum - light.ChtConfirmations) / light.ChtFrequency
|
|
||||||
}
|
|
||||||
|
|
||||||
var lastChtNum uint64
|
|
||||||
data, _ := db.Get(lastChtKey)
|
|
||||||
if len(data) == 8 {
|
|
||||||
lastChtNum = binary.BigEndian.Uint64(data[:])
|
|
||||||
}
|
|
||||||
if newChtNum <= lastChtNum {
|
|
||||||
return false
|
|
||||||
}
|
|
||||||
|
|
||||||
var t *trie.Trie
|
|
||||||
if lastChtNum > 0 {
|
|
||||||
var err error
|
|
||||||
t, err = trie.New(getChtRoot(db, lastChtNum), db)
|
|
||||||
if err != nil {
|
|
||||||
lastChtNum = 0
|
|
||||||
}
|
|
||||||
}
|
|
||||||
if lastChtNum == 0 {
|
|
||||||
t, _ = trie.New(common.Hash{}, db)
|
|
||||||
}
|
|
||||||
|
|
||||||
for num := lastChtNum * light.ChtFrequency; num < (lastChtNum+1)*light.ChtFrequency; num++ {
|
|
||||||
hash := core.GetCanonicalHash(db, num)
|
|
||||||
if hash == (common.Hash{}) {
|
|
||||||
panic("Canonical hash not found")
|
|
||||||
}
|
|
||||||
td := core.GetTd(db, hash, num)
|
|
||||||
if td == nil {
|
|
||||||
panic("TD not found")
|
|
||||||
}
|
|
||||||
var encNumber [8]byte
|
|
||||||
binary.BigEndian.PutUint64(encNumber[:], num)
|
|
||||||
var node light.ChtNode
|
|
||||||
node.Hash = hash
|
|
||||||
node.Td = td
|
|
||||||
data, _ := rlp.EncodeToBytes(node)
|
|
||||||
t.Update(encNumber[:], data)
|
|
||||||
}
|
|
||||||
|
|
||||||
root, err := t.Commit()
|
|
||||||
if err != nil {
|
|
||||||
lastChtNum = 0
|
|
||||||
} else {
|
|
||||||
lastChtNum++
|
|
||||||
|
|
||||||
log.Trace("Generated CHT", "number", lastChtNum, "root", root.Hex())
|
|
||||||
|
|
||||||
storeChtRoot(db, lastChtNum, root)
|
|
||||||
var data [8]byte
|
|
||||||
binary.BigEndian.PutUint64(data[:], lastChtNum)
|
|
||||||
db.Put(lastChtKey, data[:])
|
|
||||||
}
|
|
||||||
|
|
||||||
return newChtNum > lastChtNum
|
|
||||||
}
|
|
||||||
|
|
|
||||||
|
|
@ -95,15 +95,8 @@ func NewLightChain(odr OdrBackend, config *params.ChainConfig, engine consensus.
|
||||||
if bc.genesisBlock == nil {
|
if bc.genesisBlock == nil {
|
||||||
return nil, core.ErrNoGenesis
|
return nil, core.ErrNoGenesis
|
||||||
}
|
}
|
||||||
if bc.genesisBlock.Hash() == params.MainnetGenesisHash {
|
if ppt, ok := trustedCheckpoints[bc.genesisBlock.Hash()]; ok {
|
||||||
// add trusted CHT
|
bc.addTrustedCheckpoint(ppt)
|
||||||
WriteTrustedCht(bc.chainDb, TrustedCht{Number: 1040, Root: common.HexToHash("bb4fb4076cbe6923c8a8ce8f158452bbe19564959313466989fda095a60884ca")})
|
|
||||||
log.Info("Added trusted CHT for mainnet")
|
|
||||||
}
|
|
||||||
if bc.genesisBlock.Hash() == params.TestnetGenesisHash {
|
|
||||||
// add trusted CHT
|
|
||||||
WriteTrustedCht(bc.chainDb, TrustedCht{Number: 400, Root: common.HexToHash("2a4befa19e4675d939c3dc22dca8c6ae9fcd642be1f04b06bd6e4203cc304660")})
|
|
||||||
log.Info("Added trusted CHT for ropsten testnet")
|
|
||||||
}
|
}
|
||||||
|
|
||||||
if err := bc.loadLastState(); err != nil {
|
if err := bc.loadLastState(); err != nil {
|
||||||
|
|
@ -120,6 +113,22 @@ func NewLightChain(odr OdrBackend, config *params.ChainConfig, engine consensus.
|
||||||
return bc, nil
|
return bc, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// addTrustedCheckpoint adds a trusted checkpoint to the blockchain
|
||||||
|
func (self *LightChain) addTrustedCheckpoint(ppt trustedCheckpoint) {
|
||||||
|
if self.odr.ChtIndexer() != nil {
|
||||||
|
StoreChtRoot(self.chainDb, ppt.sectionIdx, ppt.sectionHead, ppt.chtRoot)
|
||||||
|
self.odr.ChtIndexer().AddKnownSectionHead(ppt.sectionIdx, ppt.sectionHead)
|
||||||
|
}
|
||||||
|
if self.odr.BltIndexer() != nil {
|
||||||
|
StoreBloomTrieRoot(self.chainDb, ppt.sectionIdx, ppt.sectionHead, ppt.bltRoot)
|
||||||
|
self.odr.BltIndexer().AddKnownSectionHead(ppt.sectionIdx, ppt.sectionHead)
|
||||||
|
}
|
||||||
|
if self.odr.BloomIndexer() != nil {
|
||||||
|
self.odr.BloomIndexer().AddKnownSectionHead(ppt.sectionIdx, ppt.sectionHead)
|
||||||
|
}
|
||||||
|
log.Info("Added trusted PPT", "chain name", ppt.name)
|
||||||
|
}
|
||||||
|
|
||||||
func (self *LightChain) getProcInterrupt() bool {
|
func (self *LightChain) getProcInterrupt() bool {
|
||||||
return atomic.LoadInt32(&self.procInterrupt) == 1
|
return atomic.LoadInt32(&self.procInterrupt) == 1
|
||||||
}
|
}
|
||||||
|
|
@ -449,10 +458,13 @@ func (self *LightChain) GetHeaderByNumberOdr(ctx context.Context, number uint64)
|
||||||
}
|
}
|
||||||
|
|
||||||
func (self *LightChain) SyncCht(ctx context.Context) bool {
|
func (self *LightChain) SyncCht(ctx context.Context) bool {
|
||||||
|
if self.odr.ChtIndexer() == nil {
|
||||||
|
return false
|
||||||
|
}
|
||||||
headNum := self.CurrentHeader().Number.Uint64()
|
headNum := self.CurrentHeader().Number.Uint64()
|
||||||
cht := GetTrustedCht(self.chainDb)
|
chtCount, _, _ := self.odr.ChtIndexer().Sections()
|
||||||
if headNum+1 < cht.Number*ChtFrequency {
|
if headNum+1 < chtCount*ChtFrequency {
|
||||||
num := cht.Number*ChtFrequency - 1
|
num := chtCount*ChtFrequency - 1
|
||||||
header, err := GetHeaderByNumber(ctx, self.odr, num)
|
header, err := GetHeaderByNumber(ctx, self.odr, num)
|
||||||
if header != nil && err == nil {
|
if header != nil && err == nil {
|
||||||
self.mu.Lock()
|
self.mu.Lock()
|
||||||
|
|
|
||||||
13
light/odr.go
13
light/odr.go
|
|
@ -35,6 +35,9 @@ var NoOdr = context.Background()
|
||||||
// OdrBackend is an interface to a backend service that handles ODR retrievals type
|
// OdrBackend is an interface to a backend service that handles ODR retrievals type
|
||||||
type OdrBackend interface {
|
type OdrBackend interface {
|
||||||
Database() ethdb.Database
|
Database() ethdb.Database
|
||||||
|
ChtIndexer() *core.ChainIndexer
|
||||||
|
BltIndexer() *core.ChainIndexer
|
||||||
|
BloomIndexer() *core.ChainIndexer
|
||||||
Retrieve(ctx context.Context, req OdrRequest) error
|
Retrieve(ctx context.Context, req OdrRequest) error
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -147,7 +150,8 @@ func (req *ChtRequest) StoreResult(db ethdb.Database) {
|
||||||
// BloomRequest is the ODR request type for retrieving bloom filters from a CHT structure
|
// BloomRequest is the ODR request type for retrieving bloom filters from a CHT structure
|
||||||
type BloomRequest struct {
|
type BloomRequest struct {
|
||||||
OdrRequest
|
OdrRequest
|
||||||
BltNum, BitIdx uint64
|
BltNum uint64
|
||||||
|
BitIdx uint
|
||||||
SectionIdxList []uint64
|
SectionIdxList []uint64
|
||||||
BltRoot common.Hash
|
BltRoot common.Hash
|
||||||
BloomBits [][]byte
|
BloomBits [][]byte
|
||||||
|
|
@ -157,6 +161,11 @@ type BloomRequest struct {
|
||||||
// StoreResult stores the retrieved data in local database
|
// StoreResult stores the retrieved data in local database
|
||||||
func (req *BloomRequest) StoreResult(db ethdb.Database) {
|
func (req *BloomRequest) StoreResult(db ethdb.Database) {
|
||||||
for i, sectionIdx := range req.SectionIdxList {
|
for i, sectionIdx := range req.SectionIdxList {
|
||||||
core.StoreBloomBits(db, req.BitIdx, sectionIdx, req.BloomBits[i])
|
sectionHead := core.GetCanonicalHash(db, (sectionIdx+1)*BloomTrieFrequency-1)
|
||||||
|
// if we don't have the canonical hash stored for this section head number, we'll still store it under
|
||||||
|
// a key with a zero sectionHead. GetBloomBits will look there too if we still don't have the canonical
|
||||||
|
// hash. In the unlikely case we've retrieved the section head hash since then, we'll just retrieve the
|
||||||
|
// bit vector again from the network.
|
||||||
|
core.WriteBloomBits(db, req.BitIdx, sectionIdx, sectionHead, req.BloomBits[i])
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -19,56 +19,16 @@ package light
|
||||||
import (
|
import (
|
||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
|
||||||
"math/big"
|
|
||||||
|
|
||||||
"github.com/ethereum/go-ethereum/common"
|
"github.com/ethereum/go-ethereum/common"
|
||||||
"github.com/ethereum/go-ethereum/core"
|
"github.com/ethereum/go-ethereum/core"
|
||||||
"github.com/ethereum/go-ethereum/core/types"
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
"github.com/ethereum/go-ethereum/crypto"
|
"github.com/ethereum/go-ethereum/crypto"
|
||||||
"github.com/ethereum/go-ethereum/ethdb"
|
|
||||||
"github.com/ethereum/go-ethereum/rlp"
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
)
|
)
|
||||||
|
|
||||||
var sha3_nil = crypto.Keccak256Hash(nil)
|
var sha3_nil = crypto.Keccak256Hash(nil)
|
||||||
|
|
||||||
var (
|
|
||||||
ErrNoTrustedCht = errors.New("No trusted canonical hash trie")
|
|
||||||
ErrNoHeader = errors.New("Header not found")
|
|
||||||
|
|
||||||
ChtFrequency = uint64(4096)
|
|
||||||
ChtConfirmations = uint64(2048)
|
|
||||||
trustedChtKey = []byte("TrustedCHT")
|
|
||||||
)
|
|
||||||
|
|
||||||
type ChtNode struct {
|
|
||||||
Hash common.Hash
|
|
||||||
Td *big.Int
|
|
||||||
}
|
|
||||||
|
|
||||||
type TrustedCht struct {
|
|
||||||
Number uint64
|
|
||||||
Root common.Hash
|
|
||||||
}
|
|
||||||
|
|
||||||
func GetTrustedCht(db ethdb.Database) TrustedCht {
|
|
||||||
data, _ := db.Get(trustedChtKey)
|
|
||||||
var res TrustedCht
|
|
||||||
if err := rlp.DecodeBytes(data, &res); err != nil {
|
|
||||||
return TrustedCht{0, common.Hash{}}
|
|
||||||
}
|
|
||||||
return res
|
|
||||||
}
|
|
||||||
|
|
||||||
func WriteTrustedCht(db ethdb.Database, cht TrustedCht) {
|
|
||||||
data, _ := rlp.EncodeToBytes(cht)
|
|
||||||
db.Put(trustedChtKey, data)
|
|
||||||
}
|
|
||||||
|
|
||||||
func DeleteTrustedCht(db ethdb.Database) {
|
|
||||||
db.Delete(trustedChtKey)
|
|
||||||
}
|
|
||||||
|
|
||||||
func GetHeaderByNumber(ctx context.Context, odr OdrBackend, number uint64) (*types.Header, error) {
|
func GetHeaderByNumber(ctx context.Context, odr OdrBackend, number uint64) (*types.Header, error) {
|
||||||
db := odr.Database()
|
db := odr.Database()
|
||||||
hash := core.GetCanonicalHash(db, number)
|
hash := core.GetCanonicalHash(db, number)
|
||||||
|
|
@ -81,12 +41,29 @@ func GetHeaderByNumber(ctx context.Context, odr OdrBackend, number uint64) (*typ
|
||||||
return header, nil
|
return header, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
cht := GetTrustedCht(db)
|
var (
|
||||||
if number >= cht.Number*ChtFrequency {
|
chtCount, sectionHeadNum uint64
|
||||||
|
sectionHead common.Hash
|
||||||
|
)
|
||||||
|
if odr.ChtIndexer() != nil {
|
||||||
|
chtCount, sectionHeadNum, sectionHead = odr.ChtIndexer().Sections()
|
||||||
|
canonicalHash := core.GetCanonicalHash(db, sectionHeadNum)
|
||||||
|
// if the CHT was injected as a trusted checkpoint, we have no canonical hash yet so we accept zero hash too
|
||||||
|
for chtCount > 0 && canonicalHash != sectionHead && canonicalHash != (common.Hash{}) {
|
||||||
|
chtCount--
|
||||||
|
if chtCount > 0 {
|
||||||
|
sectionHeadNum = chtCount*ChtFrequency - 1
|
||||||
|
sectionHead = odr.ChtIndexer().SectionHead(chtCount - 1)
|
||||||
|
canonicalHash = core.GetCanonicalHash(db, sectionHeadNum)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if number >= chtCount*ChtFrequency {
|
||||||
return nil, ErrNoTrustedCht
|
return nil, ErrNoTrustedCht
|
||||||
}
|
}
|
||||||
|
|
||||||
r := &ChtRequest{ChtRoot: cht.Root, ChtNum: cht.Number, BlockNum: number}
|
r := &ChtRequest{ChtRoot: GetChtRoot(db, chtCount-1, sectionHead), ChtNum: chtCount - 1, BlockNum: number}
|
||||||
if err := odr.Retrieve(ctx, r); err != nil {
|
if err := odr.Retrieve(ctx, r); err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
} else {
|
} else {
|
||||||
|
|
@ -162,3 +139,61 @@ func GetBlockReceipts(ctx context.Context, odr OdrBackend, hash common.Hash, num
|
||||||
}
|
}
|
||||||
return r.Receipts, nil
|
return r.Receipts, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// GetBloomBits retrieves a batch of compressed bloomBits vectors belonging to the given bit index and section indexes
|
||||||
|
func GetBloomBits(ctx context.Context, odr OdrBackend, bitIdx uint, sectionIdxList []uint64) ([][]byte, error) {
|
||||||
|
db := odr.Database()
|
||||||
|
result := make([][]byte, len(sectionIdxList))
|
||||||
|
var (
|
||||||
|
reqList []uint64
|
||||||
|
reqIdx []int
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
bltCount, sectionHeadNum uint64
|
||||||
|
sectionHead common.Hash
|
||||||
|
)
|
||||||
|
if odr.BltIndexer() != nil {
|
||||||
|
bltCount, sectionHeadNum, sectionHead = odr.BltIndexer().Sections()
|
||||||
|
canonicalHash := core.GetCanonicalHash(db, sectionHeadNum)
|
||||||
|
// if the BloomTrie was injected as a trusted checkpoint, we have no canonical hash yet so we accept zero hash too
|
||||||
|
for bltCount > 0 && canonicalHash != sectionHead && canonicalHash != (common.Hash{}) {
|
||||||
|
bltCount--
|
||||||
|
if bltCount > 0 {
|
||||||
|
sectionHeadNum = bltCount*BloomTrieFrequency - 1
|
||||||
|
sectionHead = odr.BltIndexer().SectionHead(bltCount - 1)
|
||||||
|
canonicalHash = core.GetCanonicalHash(db, sectionHeadNum)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
for i, sectionIdx := range sectionIdxList {
|
||||||
|
sectionHead := core.GetCanonicalHash(db, (sectionIdx+1)*BloomTrieFrequency-1)
|
||||||
|
// if we don't have the canonical hash stored for this section head number, we'll still look for
|
||||||
|
// an entry with a zero sectionHead (we store it with zero section head too if we don't know it
|
||||||
|
// at the time of the retrieval)
|
||||||
|
bloomBits, err := core.GetBloomBits(db, bitIdx, sectionIdx, sectionHead)
|
||||||
|
if err == nil {
|
||||||
|
result[i] = bloomBits
|
||||||
|
} else {
|
||||||
|
if sectionIdx >= bltCount {
|
||||||
|
return nil, ErrNoTrustedBlt
|
||||||
|
}
|
||||||
|
reqList = append(reqList, sectionIdx)
|
||||||
|
reqIdx = append(reqIdx, i)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if reqList == nil {
|
||||||
|
return result, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
r := &BloomRequest{BltRoot: GetBloomTrieRoot(db, bltCount-1, sectionHead), BltNum: bltCount - 1, BitIdx: bitIdx, SectionIdxList: reqList}
|
||||||
|
if err := odr.Retrieve(ctx, r); err != nil {
|
||||||
|
return nil, err
|
||||||
|
} else {
|
||||||
|
for i, idx := range reqIdx {
|
||||||
|
result[idx] = r.BloomBits[i]
|
||||||
|
}
|
||||||
|
return result, nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
299
light/postprocess.go
Normal file
299
light/postprocess.go
Normal file
|
|
@ -0,0 +1,299 @@
|
||||||
|
// Copyright 2016 The go-ethereum Authors
|
||||||
|
// This file is part of the go-ethereum library.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is free software: you can redistribute it and/or modify
|
||||||
|
// it under the terms of the GNU Lesser General Public License as published by
|
||||||
|
// the Free Software Foundation, either version 3 of the License, or
|
||||||
|
// (at your option) any later version.
|
||||||
|
//
|
||||||
|
// The go-ethereum library is distributed in the hope that it will be useful,
|
||||||
|
// but WITHOUT ANY WARRANTY; without even the implied warranty of
|
||||||
|
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
|
||||||
|
// GNU Lesser General Public License for more details.
|
||||||
|
//
|
||||||
|
// You should have received a copy of the GNU Lesser General Public License
|
||||||
|
// along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>.
|
||||||
|
|
||||||
|
package light
|
||||||
|
|
||||||
|
import (
|
||||||
|
"encoding/binary"
|
||||||
|
"errors"
|
||||||
|
"fmt"
|
||||||
|
"math/big"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/ethereum/go-ethereum/common"
|
||||||
|
"github.com/ethereum/go-ethereum/common/bitutil"
|
||||||
|
"github.com/ethereum/go-ethereum/core"
|
||||||
|
"github.com/ethereum/go-ethereum/core/types"
|
||||||
|
"github.com/ethereum/go-ethereum/ethdb"
|
||||||
|
"github.com/ethereum/go-ethereum/log"
|
||||||
|
"github.com/ethereum/go-ethereum/params"
|
||||||
|
"github.com/ethereum/go-ethereum/rlp"
|
||||||
|
"github.com/ethereum/go-ethereum/trie"
|
||||||
|
)
|
||||||
|
|
||||||
|
const (
|
||||||
|
ChtFrequency = 32768
|
||||||
|
ChtV1Frequency = 4096 // as long as we want to retain LES/1 compatibility, servers generate CHTs with the old, higher frequency
|
||||||
|
PPTConfirmations = 2048 // number of confirmations before a server is expected to have the given PPT available
|
||||||
|
PPTProcessConfirmations = 256 // number of confirmations before a PPT is generated
|
||||||
|
)
|
||||||
|
|
||||||
|
// trustedCheckpoint represents a set of post-processed trie roots (CHT and BloomTrie) associated with
|
||||||
|
// the appropriate section index and head hash. It is used to start light syncing from this checkpoint
|
||||||
|
// and avoid downloading the entire header chain while still being able to securely access old headers/logs.
|
||||||
|
type trustedCheckpoint struct {
|
||||||
|
name string
|
||||||
|
sectionIdx uint64
|
||||||
|
sectionHead, chtRoot, bltRoot common.Hash
|
||||||
|
}
|
||||||
|
|
||||||
|
var (
|
||||||
|
mainnetCheckpoint = trustedCheckpoint{
|
||||||
|
name: "ETH mainnet",
|
||||||
|
sectionIdx: 129,
|
||||||
|
sectionHead: common.HexToHash("64100587c8ec9a76870056d07cb0f58622552d16de6253a59cac4b580c899501"),
|
||||||
|
chtRoot: common.HexToHash("bb4fb4076cbe6923c8a8ce8f158452bbe19564959313466989fda095a60884ca"),
|
||||||
|
bltRoot: common.HexToHash("0db524b2c4a2a9520a42fd842b02d2e8fb58ff37c75cf57bd0eb82daeace6716"),
|
||||||
|
}
|
||||||
|
|
||||||
|
ropstenCheckpoint = trustedCheckpoint{
|
||||||
|
name: "Ropsten testnet",
|
||||||
|
sectionIdx: 50,
|
||||||
|
sectionHead: common.HexToHash("00bd65923a1aa67f85e6b4ae67835784dd54be165c37f056691723c55bf016bd"),
|
||||||
|
chtRoot: common.HexToHash("6f56dc61936752cc1f8c84b4addabdbe6a1c19693de3f21cb818362df2117f03"),
|
||||||
|
bltRoot: common.HexToHash("aca7d7c504d22737242effc3fdc604a762a0af9ced898036b5986c3a15220208"),
|
||||||
|
}
|
||||||
|
)
|
||||||
|
|
||||||
|
// trustedCheckpoints associates each known checkpoint with the genesis hash of the chain it belongs to
|
||||||
|
var trustedCheckpoints = map[common.Hash]trustedCheckpoint{
|
||||||
|
params.MainnetGenesisHash: mainnetCheckpoint,
|
||||||
|
params.TestnetGenesisHash: ropstenCheckpoint,
|
||||||
|
}
|
||||||
|
|
||||||
|
var (
|
||||||
|
ErrNoTrustedCht = errors.New("No trusted canonical hash trie")
|
||||||
|
ErrNoTrustedBlt = errors.New("No trusted bloom trie")
|
||||||
|
ErrNoHeader = errors.New("Header not found")
|
||||||
|
chtPrefix = []byte("chtRoot-") // chtPrefix + chtNum (uint64 big endian) -> trie root hash
|
||||||
|
ChtTablePrefix = "cht-"
|
||||||
|
)
|
||||||
|
|
||||||
|
// ChtNode structures are stored in the Canonical Hash Trie in an RLP encoded format
|
||||||
|
type ChtNode struct {
|
||||||
|
Hash common.Hash
|
||||||
|
Td *big.Int
|
||||||
|
}
|
||||||
|
|
||||||
|
// GetChtRoot reads the CHT root assoctiated to the given section from the database
|
||||||
|
// Note that sectionIdx is specified according to LES/1 CHT section size
|
||||||
|
func GetChtRoot(db ethdb.Database, sectionIdx uint64, sectionHead common.Hash) common.Hash {
|
||||||
|
var encNumber [8]byte
|
||||||
|
binary.BigEndian.PutUint64(encNumber[:], sectionIdx)
|
||||||
|
data, _ := db.Get(append(append(chtPrefix, encNumber[:]...), sectionHead.Bytes()...))
|
||||||
|
return common.BytesToHash(data)
|
||||||
|
}
|
||||||
|
|
||||||
|
// GetChtV2Root reads the CHT root assoctiated to the given section from the database
|
||||||
|
// Note that sectionIdx is specified according to LES/2 CHT section size
|
||||||
|
func GetChtV2Root(db ethdb.Database, sectionIdx uint64, sectionHead common.Hash) common.Hash {
|
||||||
|
return GetChtRoot(db, (sectionIdx+1)*(ChtFrequency/ChtV1Frequency)-1, sectionHead)
|
||||||
|
}
|
||||||
|
|
||||||
|
// StoreChtRoot writes the CHT root assoctiated to the given section into the database
|
||||||
|
// Note that sectionIdx is specified according to LES/1 CHT section size
|
||||||
|
func StoreChtRoot(db ethdb.Database, sectionIdx uint64, sectionHead, root common.Hash) {
|
||||||
|
var encNumber [8]byte
|
||||||
|
binary.BigEndian.PutUint64(encNumber[:], sectionIdx)
|
||||||
|
db.Put(append(append(chtPrefix, encNumber[:]...), sectionHead.Bytes()...), root.Bytes())
|
||||||
|
}
|
||||||
|
|
||||||
|
// ChtIndexerBackend implements core.ChainIndexerBackend
|
||||||
|
type ChtIndexerBackend struct {
|
||||||
|
db, cdb ethdb.Database
|
||||||
|
section, sectionSize uint64
|
||||||
|
lastHash common.Hash
|
||||||
|
trie *trie.Trie
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewBloomTrieIndexer creates a BloomTrie chain indexer
|
||||||
|
func NewChtIndexer(db ethdb.Database, clientMode bool) *core.ChainIndexer {
|
||||||
|
cdb := ethdb.NewTable(db, ChtTablePrefix)
|
||||||
|
idb := ethdb.NewTable(db, "chtIndex-")
|
||||||
|
var sectionSize, confirmReq uint64
|
||||||
|
if clientMode {
|
||||||
|
sectionSize = ChtFrequency
|
||||||
|
confirmReq = PPTConfirmations
|
||||||
|
} else {
|
||||||
|
sectionSize = ChtV1Frequency
|
||||||
|
confirmReq = PPTProcessConfirmations
|
||||||
|
}
|
||||||
|
return core.NewChainIndexer(db, idb, &ChtIndexerBackend{db: db, cdb: cdb, sectionSize: sectionSize}, sectionSize, confirmReq, time.Millisecond*100, "cht")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Reset implements core.ChainIndexerBackend
|
||||||
|
func (c *ChtIndexerBackend) Reset(section uint64, lastSectionHead common.Hash) {
|
||||||
|
var root common.Hash
|
||||||
|
if section > 0 {
|
||||||
|
root = GetChtRoot(c.db, section-1, lastSectionHead)
|
||||||
|
}
|
||||||
|
var err error
|
||||||
|
c.trie, err = trie.New(root, c.cdb)
|
||||||
|
if err != nil {
|
||||||
|
panic(err)
|
||||||
|
}
|
||||||
|
c.section = section
|
||||||
|
}
|
||||||
|
|
||||||
|
// Process implements core.ChainIndexerBackend
|
||||||
|
func (c *ChtIndexerBackend) Process(header *types.Header) {
|
||||||
|
hash, num := header.Hash(), header.Number.Uint64()
|
||||||
|
c.lastHash = hash
|
||||||
|
|
||||||
|
td := core.GetTd(c.db, hash, num)
|
||||||
|
if td == nil {
|
||||||
|
panic(nil)
|
||||||
|
}
|
||||||
|
var encNumber [8]byte
|
||||||
|
binary.BigEndian.PutUint64(encNumber[:], num)
|
||||||
|
data, _ := rlp.EncodeToBytes(ChtNode{hash, td})
|
||||||
|
c.trie.Update(encNumber[:], data)
|
||||||
|
}
|
||||||
|
|
||||||
|
// Commit implements core.ChainIndexerBackend
|
||||||
|
func (c *ChtIndexerBackend) Commit() error {
|
||||||
|
batch := c.cdb.NewBatch()
|
||||||
|
root, err := c.trie.CommitTo(batch)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
} else {
|
||||||
|
batch.Write()
|
||||||
|
if ((c.section+1)*c.sectionSize)%ChtFrequency == 0 {
|
||||||
|
log.Info("Storing CHT", "idx", c.section*c.sectionSize/ChtFrequency, "sectionHead", fmt.Sprintf("%064x", c.lastHash), "root", fmt.Sprintf("%064x", root))
|
||||||
|
}
|
||||||
|
StoreChtRoot(c.db, c.section, c.lastHash, root)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
const (
|
||||||
|
BloomTrieFrequency = 32768
|
||||||
|
ethBloomBitsSection = 4096
|
||||||
|
ethBloomBitsConfirmations = 256
|
||||||
|
)
|
||||||
|
|
||||||
|
var (
|
||||||
|
bloomTriePrefix = []byte("bltRoot-") // bloomTriePrefix + bloomTrieNum (uint64 big endian) -> trie root hash
|
||||||
|
BloomTrieTablePrefix = "blt-"
|
||||||
|
)
|
||||||
|
|
||||||
|
// GetBloomTrieRoot reads the BloomTrie root assoctiated to the given section from the database
|
||||||
|
func GetBloomTrieRoot(db ethdb.Database, sectionIdx uint64, sectionHead common.Hash) common.Hash {
|
||||||
|
var encNumber [8]byte
|
||||||
|
binary.BigEndian.PutUint64(encNumber[:], sectionIdx)
|
||||||
|
data, _ := db.Get(append(append(bloomTriePrefix, encNumber[:]...), sectionHead.Bytes()...))
|
||||||
|
return common.BytesToHash(data)
|
||||||
|
}
|
||||||
|
|
||||||
|
// StoreBloomTrieRoot writes the BloomTrie root assoctiated to the given section into the database
|
||||||
|
func StoreBloomTrieRoot(db ethdb.Database, sectionIdx uint64, sectionHead, root common.Hash) {
|
||||||
|
var encNumber [8]byte
|
||||||
|
binary.BigEndian.PutUint64(encNumber[:], sectionIdx)
|
||||||
|
db.Put(append(append(bloomTriePrefix, encNumber[:]...), sectionHead.Bytes()...), root.Bytes())
|
||||||
|
}
|
||||||
|
|
||||||
|
// BloomTrieIndexerBackend implements core.ChainIndexerBackend
|
||||||
|
type BloomTrieIndexerBackend struct {
|
||||||
|
db, cdb ethdb.Database
|
||||||
|
section, parentSectionSize, bloomTrieRatio uint64
|
||||||
|
trie *trie.Trie
|
||||||
|
sectionHeads []common.Hash
|
||||||
|
}
|
||||||
|
|
||||||
|
// NewBloomTrieIndexer creates a BloomTrie chain indexer
|
||||||
|
func NewBloomTrieIndexer(db ethdb.Database, clientMode bool) *core.ChainIndexer {
|
||||||
|
cdb := ethdb.NewTable(db, BloomTrieTablePrefix)
|
||||||
|
idb := ethdb.NewTable(db, "bltIndex-")
|
||||||
|
backend := &BloomTrieIndexerBackend{db: db, cdb: cdb}
|
||||||
|
var confirmReq uint64
|
||||||
|
if clientMode {
|
||||||
|
backend.parentSectionSize = BloomTrieFrequency
|
||||||
|
confirmReq = PPTConfirmations
|
||||||
|
} else {
|
||||||
|
backend.parentSectionSize = ethBloomBitsSection
|
||||||
|
confirmReq = PPTProcessConfirmations
|
||||||
|
}
|
||||||
|
backend.bloomTrieRatio = BloomTrieFrequency / backend.parentSectionSize
|
||||||
|
backend.sectionHeads = make([]common.Hash, backend.bloomTrieRatio)
|
||||||
|
return core.NewChainIndexer(db, idb, backend, BloomTrieFrequency, confirmReq-ethBloomBitsConfirmations, time.Millisecond*100, "bloomtrie")
|
||||||
|
}
|
||||||
|
|
||||||
|
// Reset implements core.ChainIndexerBackend
|
||||||
|
func (b *BloomTrieIndexerBackend) Reset(section uint64, lastSectionHead common.Hash) {
|
||||||
|
var root common.Hash
|
||||||
|
if section > 0 {
|
||||||
|
root = GetBloomTrieRoot(b.db, section-1, lastSectionHead)
|
||||||
|
}
|
||||||
|
var err error
|
||||||
|
b.trie, err = trie.New(root, b.cdb)
|
||||||
|
if err != nil {
|
||||||
|
panic(err)
|
||||||
|
}
|
||||||
|
b.section = section
|
||||||
|
}
|
||||||
|
|
||||||
|
// Process implements core.ChainIndexerBackend
|
||||||
|
func (b *BloomTrieIndexerBackend) Process(header *types.Header) {
|
||||||
|
num := header.Number.Uint64() - b.section*BloomTrieFrequency
|
||||||
|
if (num+1)%b.parentSectionSize == 0 {
|
||||||
|
b.sectionHeads[num/b.parentSectionSize] = header.Hash()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Commit implements core.ChainIndexerBackend
|
||||||
|
func (b *BloomTrieIndexerBackend) Commit() error {
|
||||||
|
var compSize, decompSize uint64
|
||||||
|
|
||||||
|
for i := uint(0); i < types.BloomBitLength; i++ {
|
||||||
|
var encKey [10]byte
|
||||||
|
binary.BigEndian.PutUint16(encKey[0:2], uint16(i))
|
||||||
|
binary.BigEndian.PutUint64(encKey[2:10], b.section)
|
||||||
|
var decomp []byte
|
||||||
|
for j := uint64(0); j < b.bloomTrieRatio; j++ {
|
||||||
|
data, err := core.GetBloomBits(b.db, i, b.section*b.bloomTrieRatio+j, b.sectionHeads[j])
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
decompData, err2 := bitutil.DecompressBytes(data, int(b.parentSectionSize/8))
|
||||||
|
if err2 != nil {
|
||||||
|
return err2
|
||||||
|
}
|
||||||
|
decomp = append(decomp, decompData...)
|
||||||
|
}
|
||||||
|
comp := bitutil.CompressBytes(decomp)
|
||||||
|
|
||||||
|
decompSize += uint64(len(decomp))
|
||||||
|
compSize += uint64(len(comp))
|
||||||
|
if len(comp) > 0 {
|
||||||
|
b.trie.Update(encKey[:], comp)
|
||||||
|
} else {
|
||||||
|
b.trie.Delete(encKey[:])
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
batch := b.cdb.NewBatch()
|
||||||
|
root, err := b.trie.CommitTo(batch)
|
||||||
|
if err != nil {
|
||||||
|
return err
|
||||||
|
} else {
|
||||||
|
batch.Write()
|
||||||
|
sectionHead := b.sectionHeads[b.bloomTrieRatio-1]
|
||||||
|
log.Info("Storing BloomTrie", "section", b.section, "sectionHead", fmt.Sprintf("%064x", sectionHead), "root", fmt.Sprintf("%064x", root), "compression ratio", float64(compSize)/float64(decompSize))
|
||||||
|
StoreBloomTrieRoot(b.db, b.section, sectionHead, root)
|
||||||
|
}
|
||||||
|
|
||||||
|
return nil
|
||||||
|
}
|
||||||
Loading…
Reference in a new issue