From 27907a59abd9dad6fa012705c94b7d1ed99043ea Mon Sep 17 00:00:00 2001 From: Arpit Temani Date: Fri, 18 Aug 2023 14:12:49 +0530 Subject: [PATCH] fix go.mod --- go.mod | 10 +- go.sum | 8 ++ les/client_handler.go | 215 +----------------------------------------- les/commons.go | 68 ------------- les/test_helper.go | 123 +++--------------------- 5 files changed, 33 insertions(+), 391 deletions(-) diff --git a/go.mod b/go.mod index 1d022174d6..000e46fb2a 100644 --- a/go.mod +++ b/go.mod @@ -8,6 +8,9 @@ require ( github.com/JekaMas/go-grpc-net-conn v0.0.0-20220708155319-6aff21f2d13d github.com/JekaMas/workerpool v1.1.8 github.com/VictoriaMetrics/fastcache v1.6.0 + github.com/aws/aws-sdk-go-v2 v1.2.0 + github.com/aws/aws-sdk-go-v2/config v1.1.1 + github.com/aws/aws-sdk-go-v2/credentials v1.1.1 github.com/aws/aws-sdk-go-v2/service/route53 v1.1.1 github.com/btcsuite/btcd/btcec/v2 v2.2.0 github.com/cespare/cp v1.1.1 @@ -54,6 +57,7 @@ require ( github.com/jedisct1/go-minisign v0.0.0-20190909160543-45766022959e github.com/julienschmidt/httprouter v1.3.0 github.com/karalabe/usb v0.0.3-0.20230711191512-61db3e06439c + github.com/kylelemons/godebug v1.1.0 github.com/maticnetwork/crand v1.0.2 github.com/maticnetwork/heimdall v0.3.1-0.20230105132832-d0063f71e3f0 github.com/maticnetwork/polyproto v0.0.2 @@ -83,6 +87,7 @@ require ( golang.org/x/text v0.9.0 golang.org/x/time v0.3.0 golang.org/x/tools v0.9.1 + gonum.org/v1/gonum v0.11.0 gopkg.in/natefinch/lumberjack.v2 v2.0.0 gopkg.in/natefinch/npipe.v2 v2.0.0-20160621034901-c1b8fa8bdcce gopkg.in/yaml.v3 v3.0.1 @@ -151,6 +156,10 @@ require ( ) require ( + github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.0.2 // indirect + github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.0.2 // indirect + github.com/aws/aws-sdk-go-v2/service/sso v1.1.1 // indirect + github.com/aws/aws-sdk-go-v2/service/sts v1.1.1 // indirect github.com/bartekn/go-bip39 v0.0.0-20171116152956-a05967ea095d // indirect github.com/bgentry/speakeasy v0.1.0 // indirect github.com/bradfitz/gomemcache v0.0.0-20190913173617-a41fca850d0b // indirect @@ -182,7 +191,6 @@ require ( github.com/json-iterator/go v1.1.12 // indirect github.com/jstemmer/go-junit-report v0.9.1 // indirect github.com/kelseyhightower/envconfig v1.4.0 // indirect - github.com/kylelemons/godebug v1.1.0 // indirect github.com/libp2p/go-buffer-pool v0.0.2 // indirect github.com/magiconair/properties v1.8.1 // indirect github.com/mitchellh/copystructure v1.0.0 // indirect diff --git a/go.sum b/go.sum index a45d0bce27..31aab0a21f 100644 --- a/go.sum +++ b/go.sum @@ -125,14 +125,21 @@ github.com/aws/aws-sdk-go v1.29.15/go.mod h1:1KvfttTE3SPKMpo8g2c6jL3ZKfXtFvKscTg github.com/aws/aws-sdk-go v1.34.28 h1:sscPpn/Ns3i0F4HPEWAVcwdIRaZZCuL7llJ2/60yPIk= github.com/aws/aws-sdk-go v1.34.28/go.mod h1:H7NKnBqNVzoTJpGfLrQkkD+ytBA93eiDYi/+8rV9s48= github.com/aws/aws-sdk-go-v2 v0.18.0/go.mod h1:JWVYvqSMppoMJC0x5wdwiImzgXTI9FuZwxzkQq9wy+g= +github.com/aws/aws-sdk-go-v2 v1.2.0 h1:BS+UYpbsElC82gB+2E2jiCBg36i8HlubTB/dO/moQ9c= github.com/aws/aws-sdk-go-v2 v1.2.0/go.mod h1:zEQs02YRBw1DjK0PoJv3ygDYOFTre1ejlJWl8FwAuQo= +github.com/aws/aws-sdk-go-v2/config v1.1.1 h1:ZAoq32boMzcaTW9bcUacBswAmHTbvlvDJICgHFZuECo= github.com/aws/aws-sdk-go-v2/config v1.1.1/go.mod h1:0XsVy9lBI/BCXm+2Tuvt39YmdHwS5unDQmxZOYe8F5Y= +github.com/aws/aws-sdk-go-v2/credentials v1.1.1 h1:NbvWIM1Mx6sNPTxowHgS2ewXCRp+NGTzUYb/96FZJbY= github.com/aws/aws-sdk-go-v2/credentials v1.1.1/go.mod h1:mM2iIjwl7LULWtS6JCACyInboHirisUUdkBPoTHMOUo= +github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.0.2 h1:EtEU7WRaWliitZh2nmuxEXrN0Cb8EgPUFGIoTMeqbzI= github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.0.2/go.mod h1:3hGg3PpiEjHnrkrlasTfxFqUsZ2GCk/fMUn4CbKgSkM= +github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.0.2 h1:4AH9fFjUlVktQMznF+YN33aWNXaR4VgDXyP28qokJC0= github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.0.2/go.mod h1:45MfaXZ0cNbeuT0KQ1XJylq8A6+OpVV2E5kvY/Kq+u8= github.com/aws/aws-sdk-go-v2/service/route53 v1.1.1 h1:cKr6St+CtC3/dl/rEBJvlk7A/IN5D5F02GNkGzfbtVU= github.com/aws/aws-sdk-go-v2/service/route53 v1.1.1/go.mod h1:rLiOUrPLW/Er5kRcQ7NkwbjlijluLsrIbu/iyl35RO4= +github.com/aws/aws-sdk-go-v2/service/sso v1.1.1 h1:37QubsarExl5ZuCBlnRP+7l1tNwZPBSTqpTBrPH98RU= github.com/aws/aws-sdk-go-v2/service/sso v1.1.1/go.mod h1:SuZJxklHxLAXgLTc1iFXbEWkXs7QRTQpCLGaKIprQW0= +github.com/aws/aws-sdk-go-v2/service/sts v1.1.1 h1:TJoIfnIFubCX0ACVeJ0w46HEH5MwjwYN4iFhuYIhfIY= github.com/aws/aws-sdk-go-v2/service/sts v1.1.1/go.mod h1:Wi0EBZwiz/K44YliU0EKxqTCJGUfYTWXrrBwkq736bM= github.com/aws/smithy-go v1.1.0 h1:D6CSsM3gdxaGaqXnPgOBCeL6Mophqzu7KJOu7zW78sU= github.com/aws/smithy-go v1.1.0/go.mod h1:EzMw8dbp/YJL4A5/sbhGddag+NPT7q084agLbB9LgIw= @@ -1546,6 +1553,7 @@ gonum.org/v1/gonum v0.0.0-20181121035319-3f7ecaa7e8ca/go.mod h1:Y+Yx5eoAFn32cQvJ gonum.org/v1/gonum v0.6.0/go.mod h1:9mxDZsDKxgMAuccQkewq682L+0eCu4dCN2yonUJTCLU= gonum.org/v1/gonum v0.8.2/go.mod h1:oe/vMfY3deqTw+1EZJhuvEW2iwGF1bW9wwu7XCu0+v0= gonum.org/v1/gonum v0.9.3/go.mod h1:TZumC3NeyVQskjXqmyWt4S3bINhy7B4eYwW69EbyX+0= +gonum.org/v1/gonum v0.11.0 h1:f1IJhK4Km5tBJmaiJXtk/PkL4cdVX6J+tGiM187uT5E= gonum.org/v1/gonum v0.11.0/go.mod h1:fSG4YDCxxUZQJ7rKsQrj0gMOg00Il0Z96/qMA4bVQhA= gonum.org/v1/netlib v0.0.0-20181029234149-ec6d1f5cefe6/go.mod h1:wa6Ws7BG/ESfp6dHfk7C6KdzKA7wR7u/rKwOGE66zvw= gonum.org/v1/netlib v0.0.0-20190313105609-8cb42192e0e0/go.mod h1:wa6Ws7BG/ESfp6dHfk7C6KdzKA7wR7u/rKwOGE66zvw= diff --git a/les/client_handler.go b/les/client_handler.go index cb229cfefa..39965d4c96 100644 --- a/les/client_handler.go +++ b/les/client_handler.go @@ -17,9 +17,6 @@ package les import ( - "context" - "math/big" - "math/rand" "sync" "time" @@ -27,11 +24,8 @@ import ( "github.com/ethereum/go-ethereum/common/mclock" "github.com/ethereum/go-ethereum/core/forkid" "github.com/ethereum/go-ethereum/core/types" - "github.com/ethereum/go-ethereum/les/downloader" "github.com/ethereum/go-ethereum/light" - "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/p2p" - "github.com/ethereum/go-ethereum/params" ) // clientHandler is responsible for receiving and processing all incoming server @@ -50,27 +44,6 @@ func newClientHandler(backend *LightEthereum) *clientHandler { backend: backend, closeCh: make(chan struct{}), } - - if ulcServers != nil { - ulc, err := newULC(ulcServers, ulcFraction) - if err != nil { - log.Error("Failed to initialize ultra light client") - } - - handler.ulc = ulc - - log.Info("Enable ultra light client mode") - } - - var height uint64 - if checkpoint != nil { - height = (checkpoint.SectionIndex+1)*params.CHTFrequency - 1 - } - - handler.fetcher = newLightFetcher(backend.blockchain, backend.engine, backend.peers, handler.ulc, backend.chainDb, backend.reqDist, handler.synchronise) - handler.downloader = downloader.New(height, backend.chainDb, backend.eventMux, nil, backend.blockchain, handler.removePeer) - handler.backend.peers.subscribe((*downloaderPeerNotify)(handler)) - return handler } @@ -81,18 +54,11 @@ func (h *clientHandler) stop() { // runPeer is the p2p protocol run function for the given version. func (h *clientHandler) runPeer(version uint, p *p2p.Peer, rw p2p.MsgReadWriter) error { - trusted := false - if h.ulc != nil { - trusted = h.ulc.trusted(p.ID()) - } - - peer := newServerPeer(int(version), h.backend.config.NetworkId, trusted, p, newMeteredMsgWriter(rw, int(version))) + peer := newServerPeer(int(version), h.backend.config.NetworkId, false, p, newMeteredMsgWriter(rw, int(version))) defer peer.close() h.wg.Add(1) - defer h.wg.Done() err := h.handle(peer, false) - return err } @@ -100,7 +66,6 @@ func (h *clientHandler) handle(p *serverPeer, noInitAnnounce bool) error { if h.backend.peers.len() >= h.backend.config.LightPeers && !p.Peer.Info().Network.Trusted { return p2p.DiscTooManyPeers } - p.Log().Debug("Light Ethereum peer connected", "name", p.Name()) // Execute the LES handshake @@ -114,7 +79,6 @@ func (h *clientHandler) handle(p *serverPeer, noInitAnnounce bool) error { if nvt, err := h.backend.serverPool.RegisterNode(p.Node()); err == nil { p.setValueTracker(nvt) p.updateVtParams() - defer func() { p.setValueTracker(nil) h.backend.serverPool.UnregisterNode(p.Node()) @@ -147,7 +111,6 @@ func (h *clientHandler) handle(p *serverPeer, noInitAnnounce bool) error { if err := h.handleMsg(p); err != nil { p.Log().Debug("Light Ethereum message handling failed", "err", err) p.fcServer.DumpLogs() - return err } } @@ -161,7 +124,6 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { if err != nil { return err } - p.Log().Trace("Light Ethereum message arrived", "code", msg.Code, "bytes", msg.Size) if msg.Size > ProtocolMaxMsgSize { @@ -175,21 +137,17 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { switch { case msg.Code == AnnounceMsg: p.Log().Trace("Received announce message") - var req announceData if err := msg.Decode(&req); err != nil { return errResp(ErrDecode, "%v: %v", msg, err) } - if err := req.sanityCheck(); err != nil { return err } - update, size := req.Update.decode() if p.rejectUpdate(size) { return errResp(ErrRequestRejected, "") } - p.updateFlowControl(update) p.updateVtParams() @@ -197,16 +155,13 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { if p.announceType == announceTypeNone { return errResp(ErrUnexpectedResponse, "") } - if p.announceType == announceTypeSigned { if err := req.checkSignature(p.ID(), update); err != nil { p.Log().Trace("Invalid announcement signature", "err", err) return err } - p.Log().Trace("Valid announcement signature") } - p.Log().Trace("Announce message content", "number", req.Number, "hash", req.Hash, "td", req.Td, "reorg", req.ReorgDepth) // Update peer head information first and then notify the announcement @@ -214,52 +169,30 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { } case msg.Code == BlockHeadersMsg: p.Log().Trace("Received block header response message") - var resp struct { ReqID, BV uint64 Headers []*types.Header } - if err := msg.Decode(&resp); err != nil { return errResp(ErrDecode, "msg %v: %v", msg, err) } - - headers := resp.Headers p.fcServer.ReceivedReply(resp.ReqID, resp.BV) p.answeredRequest(resp.ReqID) - // Filter out the explicitly requested header by the retriever - if h.backend.retriever.requested(resp.ReqID) { - deliverMsg = &Msg{ - MsgType: MsgBlockHeaders, - ReqID: resp.ReqID, - Obj: resp.Headers, - } - } else { - // Filter out any explicitly requested headers, deliver the rest to the downloader - filter := len(headers) == 1 - if filter { - headers = h.fetcher.deliverHeaders(p, resp.ReqID, resp.Headers) - } - - if len(headers) != 0 || !filter { - if err := h.downloader.DeliverHeaders(p.id, headers); err != nil { - log.Debug("Failed to deliver headers", "err", err) - } - } + deliverMsg = &Msg{ + MsgType: MsgBlockHeaders, + ReqID: resp.ReqID, + Obj: resp.Headers, } case msg.Code == BlockBodiesMsg: p.Log().Trace("Received block bodies response") - var resp struct { ReqID, BV uint64 Data []*types.Body } - if err := msg.Decode(&resp); err != nil { return errResp(ErrDecode, "msg %v: %v", msg, err) } - p.fcServer.ReceivedReply(resp.ReqID, resp.BV) p.answeredRequest(resp.ReqID) deliverMsg = &Msg{ @@ -269,16 +202,13 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { } case msg.Code == CodeMsg: p.Log().Trace("Received code response") - var resp struct { ReqID, BV uint64 Data [][]byte } - if err := msg.Decode(&resp); err != nil { return errResp(ErrDecode, "msg %v: %v", msg, err) } - p.fcServer.ReceivedReply(resp.ReqID, resp.BV) p.answeredRequest(resp.ReqID) deliverMsg = &Msg{ @@ -288,16 +218,13 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { } case msg.Code == ReceiptsMsg: p.Log().Trace("Received receipts response") - var resp struct { ReqID, BV uint64 Receipts []types.Receipts } - if err := msg.Decode(&resp); err != nil { return errResp(ErrDecode, "msg %v: %v", msg, err) } - p.fcServer.ReceivedReply(resp.ReqID, resp.BV) p.answeredRequest(resp.ReqID) deliverMsg = &Msg{ @@ -307,16 +234,13 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { } case msg.Code == ProofsV2Msg: p.Log().Trace("Received les/2 proofs response") - var resp struct { ReqID, BV uint64 Data light.NodeList } - if err := msg.Decode(&resp); err != nil { return errResp(ErrDecode, "msg %v: %v", msg, err) } - p.fcServer.ReceivedReply(resp.ReqID, resp.BV) p.answeredRequest(resp.ReqID) deliverMsg = &Msg{ @@ -326,16 +250,13 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { } case msg.Code == HelperTrieProofsMsg: p.Log().Trace("Received helper trie proof response") - var resp struct { ReqID, BV uint64 Data HelperTrieResps } - if err := msg.Decode(&resp); err != nil { return errResp(ErrDecode, "msg %v: %v", msg, err) } - p.fcServer.ReceivedReply(resp.ReqID, resp.BV) p.answeredRequest(resp.ReqID) deliverMsg = &Msg{ @@ -345,16 +266,13 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { } case msg.Code == TxStatusMsg: p.Log().Trace("Received tx status response") - var resp struct { ReqID, BV uint64 Status []light.TxStatus } - if err := msg.Decode(&resp); err != nil { return errResp(ErrDecode, "msg %v: %v", msg, err) } - p.fcServer.ReceivedReply(resp.ReqID, resp.BV) p.answeredRequest(resp.ReqID) deliverMsg = &Msg{ @@ -371,7 +289,6 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { if err := msg.Decode(&bv); err != nil { return errResp(ErrDecode, "msg %v: %v", msg, err) } - p.fcServer.ResumeFreeze(bv) p.unfreeze() p.Log().Debug("Service resumed") @@ -387,127 +304,5 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { } } } - return nil } - -func (h *clientHandler) removePeer(id string) { - h.backend.peers.unregister(id) -} - -type peerConnection struct { - handler *clientHandler - peer *serverPeer -} - -func (pc *peerConnection) Head() (common.Hash, *big.Int) { - return pc.peer.HeadAndTd() -} - -func (pc *peerConnection) RequestHeadersByHash(origin common.Hash, amount int, skip int, reverse bool) error { - rq := &distReq{ - getCost: func(dp distPeer) uint64 { - peer := dp.(*serverPeer) - return peer.getRequestCost(GetBlockHeadersMsg, amount) - }, - canSend: func(dp distPeer) bool { - return dp.(*serverPeer) == pc.peer - }, - request: func(dp distPeer) func() { - reqID := rand.Uint64() - peer := dp.(*serverPeer) - cost := peer.getRequestCost(GetBlockHeadersMsg, amount) - peer.fcServer.QueuedRequest(reqID, cost) - return func() { peer.requestHeadersByHash(reqID, origin, amount, skip, reverse) } - }, - } - - _, ok := <-pc.handler.backend.reqDist.queue(rq) - if !ok { - return light.ErrNoPeers - } - - return nil -} - -func (pc *peerConnection) RequestHeadersByNumber(origin uint64, amount int, skip int, reverse bool) error { - rq := &distReq{ - getCost: func(dp distPeer) uint64 { - peer := dp.(*serverPeer) - return peer.getRequestCost(GetBlockHeadersMsg, amount) - }, - canSend: func(dp distPeer) bool { - return dp.(*serverPeer) == pc.peer - }, - request: func(dp distPeer) func() { - reqID := rand.Uint64() - peer := dp.(*serverPeer) - cost := peer.getRequestCost(GetBlockHeadersMsg, amount) - peer.fcServer.QueuedRequest(reqID, cost) - return func() { peer.requestHeadersByNumber(reqID, origin, amount, skip, reverse) } - }, - } - - _, ok := <-pc.handler.backend.reqDist.queue(rq) - if !ok { - return light.ErrNoPeers - } - - return nil -} - -// RetrieveSingleHeaderByNumber requests a single header by the specified block -// number. This function will wait the response until it's timeout or delivered. -func (pc *peerConnection) RetrieveSingleHeaderByNumber(context context.Context, number uint64) (*types.Header, error) { - reqID := rand.Uint64() - rq := &distReq{ - getCost: func(dp distPeer) uint64 { - peer := dp.(*serverPeer) - return peer.getRequestCost(GetBlockHeadersMsg, 1) - }, - canSend: func(dp distPeer) bool { - return dp.(*serverPeer) == pc.peer - }, - request: func(dp distPeer) func() { - peer := dp.(*serverPeer) - cost := peer.getRequestCost(GetBlockHeadersMsg, 1) - peer.fcServer.QueuedRequest(reqID, cost) - return func() { peer.requestHeadersByNumber(reqID, number, 1, 0, false) } - }, - } - - var header *types.Header - - if err := pc.handler.backend.retriever.retrieve(context, reqID, rq, func(peer distPeer, msg *Msg) error { - if msg.MsgType != MsgBlockHeaders { - return errInvalidMessageType - } - headers := msg.Obj.([]*types.Header) - if len(headers) != 1 { - return errInvalidEntryCount - } - header = headers[0] - return nil - }, nil); err != nil { - return nil, err - } - - return header, nil -} - -// downloaderPeerNotify implements peerSetNotify -type downloaderPeerNotify clientHandler - -func (d *downloaderPeerNotify) registerPeer(p *serverPeer) { - h := (*clientHandler)(d) - pc := &peerConnection{ - handler: h, - peer: p, - } - h.downloader.RegisterLightPeer(p.id, eth.ETH66, pc) -} - -func (d *downloaderPeerNotify) unregisterPeer(p *serverPeer) { - h := (*clientHandler)(d) - h.downloader.UnregisterPeer(p.id) -} diff --git a/les/commons.go b/les/commons.go index ca514fc8eb..cb3fc430b7 100644 --- a/les/commons.go +++ b/les/commons.go @@ -27,10 +27,7 @@ import ( "github.com/ethereum/go-ethereum/core/types" "github.com/ethereum/go-ethereum/eth/ethconfig" "github.com/ethereum/go-ethereum/ethdb" - "github.com/ethereum/go-ethereum/les/checkpointoracle" "github.com/ethereum/go-ethereum/light" - "github.com/ethereum/go-ethereum/log" - "github.com/ethereum/go-ethereum/node" "github.com/ethereum/go-ethereum/p2p" "github.com/ethereum/go-ethereum/p2p/enode" "github.com/ethereum/go-ethereum/params" @@ -71,7 +68,6 @@ type NodeInfo struct { // makeProtocols creates protocol descriptors for the given LES versions. func (c *lesCommons) makeProtocols(versions []uint, runPeer func(version uint, p *p2p.Peer, rw p2p.MsgReadWriter) error, peerInfo func(id enode.ID) interface{}, dialCandidates enode.Iterator) []p2p.Protocol { protos := make([]p2p.Protocol, len(versions)) - for i, version := range versions { version := version protos[i] = p2p.Protocol{ @@ -86,7 +82,6 @@ func (c *lesCommons) makeProtocols(versions []uint, runPeer func(version uint, p DialCandidates: dialCandidates, } } - return protos } @@ -94,7 +89,6 @@ func (c *lesCommons) makeProtocols(versions []uint, runPeer func(version uint, p func (c *lesCommons) nodeInfo() interface{} { head := c.chainReader.CurrentHeader() hash := head.Hash() - return &NodeInfo{ Network: c.config.NetworkId, Difficulty: rawdb.ReadTd(c.chainDb, hash, head.Number.Uint64()), @@ -103,65 +97,3 @@ func (c *lesCommons) nodeInfo() interface{} { Head: hash, } } - -// latestLocalCheckpoint finds the common stored section index and returns a set -// of post-processed trie roots (CHT and BloomTrie) associated with the appropriate -// section index and head hash as a local checkpoint package. -func (c *lesCommons) latestLocalCheckpoint() params.TrustedCheckpoint { - sections, _, _ := c.chtIndexer.Sections() - sections2, _, _ := c.bloomTrieIndexer.Sections() - // Cap the section index if the two sections are not consistent. - if sections > sections2 { - sections = sections2 - } - - if sections == 0 { - // No checkpoint information can be provided. - return params.TrustedCheckpoint{} - } - - return c.localCheckpoint(sections - 1) -} - -// localCheckpoint returns a set of post-processed trie roots (CHT and BloomTrie) -// associated with the appropriate head hash by specific section index. -// -// The returned checkpoint is only the checkpoint generated by the local indexers, -// not the stable checkpoint registered in the registrar contract. -func (c *lesCommons) localCheckpoint(index uint64) params.TrustedCheckpoint { - sectionHead := c.chtIndexer.SectionHead(index) - - return params.TrustedCheckpoint{ - SectionIndex: index, - SectionHead: sectionHead, - CHTRoot: light.GetChtRoot(c.chainDb, index, sectionHead), - BloomRoot: light.GetBloomTrieRoot(c.chainDb, index, sectionHead), - } -} - -// setupOracle sets up the checkpoint oracle contract client. -func (c *lesCommons) setupOracle(node *node.Node, genesis common.Hash, ethconfig *ethconfig.Config) *checkpointoracle.CheckpointOracle { - config := ethconfig.CheckpointOracle - if config == nil { - // Try loading default config. - config = params.CheckpointOracles[genesis] - } - - if config == nil { - log.Info("Checkpoint oracle is not enabled") - return nil - } - - if config.Address == (common.Address{}) || uint64(len(config.Signers)) < config.Threshold { - log.Warn("Invalid checkpoint oracle config") - return nil - } - - oracle := checkpointoracle.New(config, c.localCheckpoint) - rpcClient, _ := node.Attach() - client := ethclient.NewClient(rpcClient) - oracle.Start(client) - log.Info("Configured checkpoint oracle", "address", config.Address, "signers", len(config.Signers), "threshold", config.Threshold) - - return oracle -} diff --git a/les/test_helper.go b/les/test_helper.go index 9d59da194f..b03bca14bf 100644 --- a/les/test_helper.go +++ b/les/test_helper.go @@ -25,17 +25,14 @@ import ( "crypto/rand" "fmt" "math/big" - "sync/atomic" "testing" "time" - "github.com/ethereum/go-ethereum/accounts/abi/bind" "github.com/ethereum/go-ethereum/accounts/abi/bind/backends" "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common/mclock" "github.com/ethereum/go-ethereum/consensus" "github.com/ethereum/go-ethereum/consensus/ethash" - "github.com/ethereum/go-ethereum/contracts/checkpointoracle/contract" "github.com/ethereum/go-ethereum/core" "github.com/ethereum/go-ethereum/core/forkid" "github.com/ethereum/go-ethereum/core/rawdb" @@ -46,7 +43,6 @@ import ( "github.com/ethereum/go-ethereum/eth/ethconfig" "github.com/ethereum/go-ethereum/ethdb" "github.com/ethereum/go-ethereum/event" - "github.com/ethereum/go-ethereum/les/checkpointoracle" "github.com/ethereum/go-ethereum/les/flowcontrol" vfs "github.com/ethereum/go-ethereum/les/vflux/server" "github.com/ethereum/go-ethereum/light" @@ -106,16 +102,12 @@ func prepare(n int, backend *backends.SimulatedBackend) { ctx = context.Background() signer = types.HomesteadSigner{} ) - for i := 0; i < n; i++ { switch i { case 0: // Builtin-block // number: 1 - // txs: 2 - // deploy checkpoint contract - auth, _ := bind.NewKeyedTransactorWithChainID(bankKey, big.NewInt(1337)) - oracleAddr, _, _, _ = contract.DeployCheckpointOracle(auth, backend, []common.Address{signerAddr}, sectionSize, processConfirms, big.NewInt(1)) + // txs: 1 // bankUser transfers some ether to user1 nonce, _ := backend.PendingNonceAt(ctx, bankAddr) @@ -125,6 +117,7 @@ func prepare(n int, backend *backends.SimulatedBackend) { // Builtin-block // number: 2 // txs: 4 + bankNonce, _ := backend.PendingNonceAt(ctx, bankAddr) userNonce1, _ := backend.PendingNonceAt(ctx, userAddr1) @@ -139,7 +132,6 @@ func prepare(n int, backend *backends.SimulatedBackend) { // user1 deploys a test contract tx3, _ := types.SignTx(types.NewContractCreation(userNonce1+1, big.NewInt(0), 200000, big.NewInt(params.InitialBaseFee), testContractCode), signer, userKey1) backend.SendTransaction(ctx, tx3) - testContractAddr = crypto.CreateAddress(userAddr1, userNonce1+1) // user1 deploys a event contract @@ -149,6 +141,7 @@ func prepare(n int, backend *backends.SimulatedBackend) { // Builtin-block // number: 3 // txs: 2 + // bankUser transfer some ether to signer bankNonce, _ := backend.PendingNonceAt(ctx, bankAddr) tx1, _ := types.SignTx(types.NewTransaction(bankNonce, signerAddr, big.NewInt(1000000000), params.TxGas, big.NewInt(params.InitialBaseFee), nil), signer, bankKey) @@ -162,13 +155,13 @@ func prepare(n int, backend *backends.SimulatedBackend) { // Builtin-block // number: 4 // txs: 1 + // invoke test contract bankNonce, _ := backend.PendingNonceAt(ctx, bankAddr) data := common.Hex2Bytes("C16431B900000000000000000000000000000000000000000000000000000000000000020000000000000000000000000000000000000000000000000000000000000002") tx, _ := types.SignTx(types.NewTransaction(bankNonce, testContractAddr, big.NewInt(0), 100000, big.NewInt(params.InitialBaseFee), data), signer, bankKey) backend.SendTransaction(ctx, tx) } - backend.Commit() } } @@ -181,7 +174,6 @@ func testIndexers(db ethdb.Database, odr light.OdrBackend, config *light.Indexer indexers[2] = light.NewBloomTrieIndexer(db, odr, config.BloomSize, config.BloomTrieSize, disablePruning) // make bloomTrieIndexer as a child indexer of bloom indexer. indexers[1].AddChildIndexer(indexers[2]) - return indexers[:] } @@ -196,29 +188,8 @@ func newTestClientHandler(backend *backends.SimulatedBackend, odr *LesOdr, index BaseFee: big.NewInt(params.InitialBaseFee), } ) - genesis := gspec.MustCommit(db) - chain, _ := light.NewLightChain(odr, gspec.Config, engine, nil, nil) - - if indexers != nil { - checkpointConfig := ¶ms.CheckpointOracleConfig{ - Address: crypto.CreateAddress(bankAddr, 0), - Signers: []common.Address{signerAddr}, - Threshold: 1, - } - getLocal := func(index uint64) params.TrustedCheckpoint { - chtIndexer := indexers[0] - sectionHead := chtIndexer.SectionHead(index) - - return params.TrustedCheckpoint{ - SectionIndex: index, - SectionHead: sectionHead, - CHTRoot: light.GetChtRoot(db, index, sectionHead), - BloomRoot: light.GetBloomTrieRoot(db, index, sectionHead), - } - } - oracle = checkpointoracle.New(checkpointConfig, getLocal) - } + chain, _ := light.NewLightChain(odr, gspec.Config, engine) client := &LightEthereum{ lesCommons: lesCommons{ @@ -241,12 +212,6 @@ func newTestClientHandler(backend *backends.SimulatedBackend, odr *LesOdr, index } client.handler = newClientHandler(client) - if client.oracle != nil { - client.oracle.Start(backend) - } - - client.handler.start() - return client.handler, func() { client.handler.stop() } @@ -261,7 +226,6 @@ func newTestServerHandler(blocks int, indexers []*core.ChainIndexer, db ethdb.Da BaseFee: big.NewInt(params.InitialBaseFee), } ) - genesis := gspec.MustCommit(db) // create a simulation backend and pre-commit several customized block to the database. @@ -270,27 +234,9 @@ func newTestServerHandler(blocks int, indexers []*core.ChainIndexer, db ethdb.Da txpoolConfig := legacypool.DefaultConfig txpoolConfig.Journal = "" - txpool := txpool.NewTxPool(txpoolConfig, gspec.Config, simulation.Blockchain()) - if indexers != nil { - checkpointConfig := ¶ms.CheckpointOracleConfig{ - Address: crypto.CreateAddress(bankAddr, 0), - Signers: []common.Address{signerAddr}, - Threshold: 1, - } - getLocal := func(index uint64) params.TrustedCheckpoint { - chtIndexer := indexers[0] - sectionHead := chtIndexer.SectionHead(index) - - return params.TrustedCheckpoint{ - SectionIndex: index, - SectionHead: sectionHead, - CHTRoot: light.GetChtRoot(db, index, sectionHead), - BloomRoot: light.GetBloomTrieRoot(db, index, sectionHead), - } - } - oracle = checkpointoracle.New(checkpointConfig, getLocal) - } + pool := legacypool.New(txpoolConfig, simulation.Blockchain()) + txpool, _ := txpool.New(new(big.Int).SetUint64(txpoolConfig.PriceLimit), simulation.Blockchain(), []txpool.SubPool{pool}) server := &LesServer{ lesCommons: lesCommons{ @@ -315,17 +261,10 @@ func newTestServerHandler(blocks int, indexers []*core.ChainIndexer, db ethdb.Da server.clientPool = vfs.NewClientPool(db, testBufRecharge, defaultConnectedBias, clock, alwaysTrueFn) server.clientPool.Start() server.clientPool.SetLimits(10000, 10000) // Assign enough capacity for clientpool - server.handler = newServerHandler(server, simulation.Blockchain(), db, txpool, func() bool { return true }) - if server.oracle != nil { - server.oracle.Start(simulation) - } - server.servingQueue.setThreads(4) server.handler.start() - closer := func() { server.Stop() } - return server.handler, simulation, closer } @@ -348,7 +287,6 @@ func (p *testPeer) handshakeWithServer(t *testing.T, td *big.Int, head common.Ha if p.cpeer == nil { t.Fatal("handshake for client peer only") } - var sendList keyValueList sendList = sendList.add("protocolVersion", uint64(p.cpeer.version)) sendList = sendList.add("networkId", uint64(NetworkId)) @@ -356,15 +294,12 @@ func (p *testPeer) handshakeWithServer(t *testing.T, td *big.Int, head common.Ha sendList = sendList.add("headHash", head) sendList = sendList.add("headNum", headNum) sendList = sendList.add("genesisHash", genesis) - if p.cpeer.version >= lpv4 { sendList = sendList.add("forkID", &forkID) } - if err := p2p.ExpectMsg(p.app, StatusMsg, nil); err != nil { t.Fatalf("status recv: %v", err) } - if err := p2p.Send(p.app, StatusMsg, &sendList); err != nil { t.Fatalf("status send: %v", err) } @@ -377,7 +312,6 @@ func (p *testPeer) handshakeWithServer(t *testing.T, td *big.Int, head common.Ha if p.speer == nil { t.Fatal("handshake for server peer only") } - var sendList keyValueList sendList = sendList.add("protocolVersion", uint64(p.speer.version)) sendList = sendList.add("networkId", uint64(NetworkId)) @@ -388,21 +322,18 @@ func (p *testPeer) handshakeWithServer(t *testing.T, td *big.Int, head common.Ha sendList = sendList.add("serveHeaders", nil) sendList = sendList.add("serveChainSince", uint64(0)) sendList = sendList.add("serveStateSince", uint64(0)) - sendList = sendList.add("serveRecentState", core.DefaultCacheConfig.TriesInMemory-4) + sendList = sendList.add("serveRecentState", uint64(core.TriesInMemory-4)) sendList = sendList.add("txRelay", nil) sendList = sendList.add("flowControl/BL", testBufLimit) sendList = sendList.add("flowControl/MRR", testBufRecharge) sendList = sendList.add("flowControl/MRC", costList) - if p.speer.version >= lpv4 { sendList = sendList.add("forkID", &forkID) sendList = sendList.add("recentTxLookup", recentTxLookup) } - if err := p2p.ExpectMsg(p.app, StatusMsg, nil); err != nil { t.Fatalf("status recv: %v", err) } - if err := p2p.Send(p.app, StatusMsg, &sendList); err != nil { t.Fatalf("status send: %v", err) } @@ -420,7 +351,6 @@ func newTestPeerPair(name string, version int, server *serverHandler, client *cl // Generate a random id and create the peer var id enode.ID - rand.Read(id[:]) peer1 := newClientPeer(version, NetworkId, p2p.NewPeer(id, name, nil), net) @@ -429,7 +359,6 @@ func newTestPeerPair(name string, version int, server *serverHandler, client *cl // Start the peer on a new thread errc1 := make(chan error, 1) errc2 := make(chan error, 1) - go func() { select { case <-server.closeCh: @@ -453,14 +382,11 @@ func newTestPeerPair(name string, version int, server *serverHandler, client *cl return nil, nil, fmt.Errorf("failed to establish protocol connection %v", err) default: } - - if atomic.LoadUint32(&peer1.serving) == 1 && atomic.LoadUint32(&peer2.serving) == 1 { + if peer1.serving.Load() && peer2.serving.Load() { break } - time.Sleep(50 * time.Millisecond) } - return &testPeer{cpeer: peer1, net: net, app: app}, &testPeer{speer: peer2, net: app, app: net}, nil } @@ -486,13 +412,11 @@ type testClient struct { // Generate a random id and create the peer var id enode.ID - rand.Read(id[:]) peer := newServerPeer(version, NetworkId, false, p2p.NewPeer(id, name, nil), net) // Start the peer on a new thread errCh := make(chan error, 1) - go func() { select { case <-client.handler.closeCh: @@ -500,19 +424,16 @@ type testClient struct { case errCh <- client.handler.handle(peer, false): } }() - tp := &testPeer{ app: app, net: net, speer: peer, } - var ( genesis = client.handler.backend.blockchain.Genesis() head = client.handler.backend.blockchain.CurrentHeader() td = client.handler.backend.blockchain.GetTd(head.Hash(), head.Number.Uint64()) ) - forkID := forkid.NewID(client.handler.backend.blockchain.Config(), genesis.Hash(), head.Number.Uint64(), head.Time) tp.handshakeWithClient(t, td, head.Hash(), head.Number.Uint64(), genesis.Hash(), forkID, testCostList(0), recentTxLookup) // disable flow control by default @@ -523,19 +444,15 @@ type testClient struct { return nil, nil, nil default: } - - if atomic.LoadUint32(&peer.serving) == 1 { + if peer.serving.Load() { break } - time.Sleep(50 * time.Millisecond) } - closePeer := func() { tp.speer.close() tp.close() } - return tp, closePeer, errCh }*/ @@ -559,13 +476,11 @@ func (server *testServer) newRawPeer(t *testing.T, name string, version int) (*t // Generate a random id and create the peer var id enode.ID - rand.Read(id[:]) peer := newClientPeer(version, NetworkId, p2p.NewPeer(id, name, nil), net) // Start the peer on a new thread errCh := make(chan error, 1) - go func() { select { case <-server.handler.closeCh: @@ -573,19 +488,16 @@ func (server *testServer) newRawPeer(t *testing.T, name string, version int) (*t case errCh <- server.handler.handle(peer): } }() - tp := &testPeer{ app: app, net: net, cpeer: peer, } - var ( genesis = server.handler.blockchain.Genesis() head = server.handler.blockchain.CurrentHeader() td = server.handler.blockchain.GetTd(head.Hash(), head.Number.Uint64()) ) - forkID := forkid.NewID(server.handler.blockchain.Config(), genesis.Hash(), head.Number.Uint64(), head.Time) tp.handshakeWithServer(t, td, head.Hash(), head.Number.Uint64(), genesis.Hash(), forkID) @@ -596,19 +508,15 @@ func (server *testServer) newRawPeer(t *testing.T, name string, version int) (*t return nil, nil, nil default: } - - if atomic.LoadUint32(&peer.serving) == 1 { + if peer.serving.Load() { break } - time.Sleep(50 * time.Millisecond) } - closePeer := func() { tp.cpeer.close() tp.close() } - return tp, closePeer, errCh } @@ -628,12 +536,10 @@ func newClientServerEnv(t *testing.T, config testnetConfig) (*testServer, *testC cdb = rawdb.NewMemoryDatabase() speers = newServerPeerSet() ) - var clock mclock.Clock = &mclock.System{} if config.simClock { clock = &mclock.Simulated{} } - dist := newRequestDistributor(speers, clock) rm := newRetrieveManager(speers, dist, func() time.Duration { return time.Millisecond * 500 }) odr := NewLesOdr(cdb, light.TestClientIndexerConfig, speers, rm) @@ -656,16 +562,12 @@ func newClientServerEnv(t *testing.T, config testnetConfig) (*testServer, *testC if config.indexFn != nil { config.indexFn(scIndexer, sbIndexer, sbtIndexer) } - var ( err error speer, cpeer *testPeer ) - if config.connect { done := make(chan struct{}) - client.syncEnd = func(_ *types.Header) { close(done) } - cpeer, speer, err = newTestPeerPair("peer", config.protocol, server, client, false) if err != nil { t.Fatalf("Failed to connect testing peers %v", err) @@ -676,7 +578,6 @@ func newClientServerEnv(t *testing.T, config testnetConfig) (*testServer, *testC t.Fatal("test peer did not connect and sync within 3s") } } - s := &testServer{ clock: clock, backend: b, @@ -703,7 +604,6 @@ func newClientServerEnv(t *testing.T, config testnetConfig) (*testServer, *testC cpeer.cpeer.close() speer.speer.close() } - ccIndexer.Close() cbIndexer.Close() scIndexer.Close() @@ -713,7 +613,6 @@ func newClientServerEnv(t *testing.T, config testnetConfig) (*testServer, *testC b.Close() clientClose() } - return s, c, teardown }