From 9be4810e77c668ffcd847ec8983c948a987c8062 Mon Sep 17 00:00:00 2001 From: Zsolt Felfoldi Date: Fri, 28 Feb 2020 20:26:51 +0100 Subject: [PATCH] les: fixed compile errors --- les/api.go | 37 ++++++++++++++++++++------------- les/client.go | 2 +- les/client_handler.go | 4 ++-- les/clientpool.go | 4 ++-- les/clientpool_test.go | 2 +- les/handler_test.go | 21 ++++++++++--------- les/peer.go | 46 ++++++------------------------------------ les/peer_test.go | 2 +- les/server.go | 2 +- les/server_handler.go | 12 +++++------ les/test_helper.go | 13 +++++++----- 11 files changed, 61 insertions(+), 84 deletions(-) diff --git a/les/api.go b/les/api.go index 7590925a6c..1e689e3582 100644 --- a/les/api.go +++ b/les/api.go @@ -296,7 +296,7 @@ func (api *PrivateDebugAPI) FreezeClient(id enode.ID) error { if c == nil { return fmt.Errorf("client %064x is not connected", id[:]) } - c.peer.freezeClient() + c.peer.freeze() return nil }) } @@ -355,16 +355,18 @@ func (api *PrivateLightAPI) GetCheckpointContractAddress() (string, error) { // PrivateLespayAPI provides an API to use the LESpay commands of either the local or a remote server type PrivateLespayAPI struct { - peerSet *peerSet + clientPeerSet *clientPeerSet + serverPeerSet *serverPeerSet clientHandler *clientHandler dht *discv5.Network tokenSale *tokenSale } // NewPrivateLespayAPI creates a new LESPAY API. -func NewPrivateLespayAPI(peerSet *peerSet, clientHandler *clientHandler, dht *discv5.Network, tokenSale *tokenSale) *PrivateLespayAPI { +func NewPrivateLespayAPI(clientPeerSet *clientPeerSet, serverPeerSet *serverPeerSet, clientHandler *clientHandler, dht *discv5.Network, tokenSale *tokenSale) *PrivateLespayAPI { return &PrivateLespayAPI{ - peerSet: peerSet, + clientPeerSet: clientPeerSet, + serverPeerSet: serverPeerSet, clientHandler: clientHandler, dht: dht, tokenSale: tokenSale, @@ -379,18 +381,25 @@ func NewPrivateLespayAPI(peerSet *peerSet, clientHandler *clientHandler, dht *di // If remote is false then the command is executed locally, with the specified remote node assumed as sender. func (api *PrivateLespayAPI) makeCall(ctx context.Context, remote bool, nodeStr string, cmd []byte) ([]byte, error) { var ( - id enode.ID - freeID string - peer *peer - node *enode.Node - err error + id enode.ID + freeID string + clientPeer *clientPeer + serverPeer *serverPeer + node *enode.Node + err error ) if nodeStr != "" { if id, err = enode.ParseID(nodeStr); err == nil { - if peer = api.peerSet.Peer(peerIdToString(id)); peer == nil { - return nil, errors.New("peer not connected") + if api.clientPeerSet != nil { + if clientPeer = api.clientPeerSet.peer(peerIdToString(id)); clientPeer == nil { + return nil, errors.New("peer not connected") + } + freeID = clientPeer.freeClientId() + } else { + if serverPeer = api.serverPeerSet.peer(peerIdToString(id)); serverPeer == nil { + return nil, errors.New("peer not connected") + } } - freeID = peer.freeClientId() } else { var err error if node, err = enode.Parse(enode.ValidSchemes, nodeStr); err == nil { @@ -408,12 +417,12 @@ func (api *PrivateLespayAPI) makeCall(ctx context.Context, remote bool, nodeStr cancelFn func() bool ) delivered := make(chan struct{}) - if peer != nil { + if serverPeer != nil { // remote call to a connected peer through LES if api.clientHandler == nil { return nil, errors.New("client handler not available") } - cancelFn = api.clientHandler.makeLespayCall(peer, cmd, func(r []byte, delay uint) bool { + cancelFn = api.clientHandler.makeLespayCall(serverPeer, cmd, func(r []byte, delay uint) bool { reply = r close(delivered) return reply != nil diff --git a/les/client.go b/les/client.go index c1c5b354c8..643aa88dc4 100644 --- a/les/client.go +++ b/les/client.go @@ -212,7 +212,7 @@ func (s *LightEthereum) APIs() []rpc.API { { Namespace: "lespay", Version: "1.0", - Service: NewPrivateLespayAPI(s.lesCommons.peers, s.handler, s.srvr.DiscV5, nil), + Service: NewPrivateLespayAPI(nil, s.peers, s.handler, s.srvr.DiscV5, nil), Public: false, }, }...) diff --git a/les/client_handler.go b/les/client_handler.go index a54187cb8e..56392632bc 100644 --- a/les/client_handler.go +++ b/les/client_handler.go @@ -391,12 +391,12 @@ func (h *clientHandler) handleMsg(p *serverPeer) error { // makeLespayCall sends a lespay command through an LES connection and registers // a response handler. It returns a cancel function that removes the response // handler and calls it with a nil parameter if the response has not arrived yet. -func (h *clientHandler) makeLespayCall(p *peer, cmd []byte, handler func([]byte, uint) bool) func() bool { +func (h *clientHandler) makeLespayCall(p *serverPeer, cmd []byte, handler func([]byte, uint) bool) func() bool { reqID := genReqID() h.lespayReplyLock.Lock() h.lespayReplyHandlers[reqID] = handler h.lespayReplyLock.Unlock() - if p.SendLespay(reqID, cmd) != nil { + if p.sendLespay(reqID, cmd) != nil { h.lespayReplyLock.Lock() delete(h.lespayReplyHandlers, reqID) h.lespayReplyLock.Unlock() diff --git a/les/clientpool.go b/les/clientpool.go index 760c24d1e5..80915317ca 100644 --- a/les/clientpool.go +++ b/les/clientpool.go @@ -123,7 +123,7 @@ type clientPoolPeer interface { ID() enode.ID freeClientId() string updateCapacity(uint64) - freezeClient() + freeze() } // clientInfo represents a connected client @@ -500,7 +500,7 @@ func (f *clientPool) disconnect(p clientPoolPeer) { } // drop deactivates the peer if necessary and drops it from the inactive queue -func (f *clientPool) drop(p clientPeer, kicked bool) { +func (f *clientPool) drop(p clientPoolPeer, kicked bool) { // Short circuit if client pool is already closed. if f.closed { return diff --git a/les/clientpool_test.go b/les/clientpool_test.go index dd0818c3a7..94ccb33209 100644 --- a/les/clientpool_test.go +++ b/les/clientpool_test.go @@ -81,7 +81,7 @@ func (i *poolTestPeer) updateCapacity(cap uint64) { } } -func (i *poolTestPeer) freezeClient() {} +func (i *poolTestPeer) freeze() {} func testClientPool(t *testing.T, activeLimit, clientCount, paidCount int, randomDisconnect bool) { rand.Seed(time.Now().UnixNano()) diff --git a/les/handler_test.go b/les/handler_test.go index 8107270354..d3e4df1b77 100644 --- a/les/handler_test.go +++ b/les/handler_test.go @@ -177,7 +177,7 @@ func testGetBlockHeaders(t *testing.T, protocol int) { // Send the hash request and verify the response reqID++ - cost := server.peer.peer.GetRequestCost(GetBlockHeadersMsg, int(tt.query.Amount)) + cost := server.peer.speer.getRequestCost(GetBlockHeadersMsg, int(tt.query.Amount)) sendRequest(server.peer.app, GetBlockHeadersMsg, reqID, tt.query) if err := expectResponse(server.peer.app, protocol, BlockHeadersMsg, reqID, testBufLimit, cost, headers); err != nil { t.Errorf("test %d: headers mismatch: %v", i, err) @@ -255,7 +255,7 @@ func testGetBlockBodies(t *testing.T, protocol int) { reqID++ // Send the hash request and verify the response - cost := server.peer.peer.GetRequestCost(GetBlockBodiesMsg, len(hashes)) + cost := server.peer.speer.getRequestCost(GetBlockBodiesMsg, len(hashes)) sendRequest(server.peer.app, GetBlockBodiesMsg, reqID, hashes) if err := expectResponse(server.peer.app, protocol, BlockBodiesMsg, reqID, testBufLimit, cost, bodies); err != nil { t.Errorf("test %d: bodies mismatch: %v", i, err) @@ -287,7 +287,7 @@ func testGetCode(t *testing.T, protocol int) { } } - cost := server.peer.peer.GetRequestCost(GetCodeMsg, len(codereqs)) + cost := server.peer.speer.getRequestCost(GetCodeMsg, len(codereqs)) sendRequest(server.peer.app, GetCodeMsg, 42, codereqs) if err := expectResponse(server.peer.app, protocol, CodeMsg, 42, testBufLimit, cost, codes); err != nil { t.Errorf("codes mismatch: %v", err) @@ -308,7 +308,7 @@ func testGetStaleCode(t *testing.T, protocol int) { BHash: bc.GetHeaderByNumber(number).Hash(), AccKey: crypto.Keccak256(testContractAddr[:]), } - cost := server.peer.peer.GetRequestCost(GetCodeMsg, 1) + cost := server.peer.speer.getRequestCost(GetCodeMsg, 1) sendRequest(server.peer.app, GetCodeMsg, 42, []*CodeReq{req}) if err := expectResponse(server.peer.app, protocol, CodeMsg, 42, testBufLimit, cost, expected); err != nil { t.Errorf("codes mismatch: %v", err) @@ -340,7 +340,7 @@ func testGetReceipt(t *testing.T, protocol int) { receipts = append(receipts, rawdb.ReadRawReceipts(server.db, block.Hash(), block.NumberU64())) } // Send the hash request and verify the response - cost := server.peer.peer.GetRequestCost(GetReceiptsMsg, len(hashes)) + cost := server.peer.speer.getRequestCost(GetReceiptsMsg, len(hashes)) sendRequest(server.peer.app, GetReceiptsMsg, 42, hashes) if err := expectResponse(server.peer.app, protocol, ReceiptsMsg, 42, testBufLimit, cost, receipts); err != nil { t.Errorf("receipts mismatch: %v", err) @@ -376,7 +376,7 @@ func testGetProofs(t *testing.T, protocol int) { } } // Send the proof request and verify the response - cost := server.peer.peer.GetRequestCost(GetProofsV2Msg, len(proofreqs)) + cost := server.peer.speer.getRequestCost(GetProofsV2Msg, len(proofreqs)) sendRequest(server.peer.app, GetProofsV2Msg, 42, proofreqs) if err := expectResponse(server.peer.app, protocol, ProofsV2Msg, 42, testBufLimit, cost, proofsV2.NodeList()); err != nil { t.Errorf("proofs mismatch: %v", err) @@ -410,6 +410,7 @@ func testGetStaleProof(t *testing.T, protocol int) { t.Prove(account, 0, proofsV2) expected = proofsV2.NodeList() } + cost := server.peer.speer.getRequestCost(GetProofsV2Msg, 1) if err := expectResponse(server.peer.app, protocol, ProofsV2Msg, 42, testBufLimit, cost, expected); err != nil { t.Errorf("codes mismatch: %v", err) } @@ -461,7 +462,7 @@ func testGetCHTProofs(t *testing.T, protocol int) { AuxReq: auxHeader, }} // Send the proof request and verify the response - cost := server.peer.peer.GetRequestCost(GetHelperTrieProofsMsg, len(requestsV2)) + cost := server.peer.speer.getRequestCost(GetHelperTrieProofsMsg, len(requestsV2)) sendRequest(server.peer.app, GetHelperTrieProofsMsg, 42, requestsV2) if err := expectResponse(server.peer.app, protocol, HelperTrieProofsMsg, 42, testBufLimit, cost, proofsV2); err != nil { t.Errorf("proofs mismatch: %v", err) @@ -510,7 +511,7 @@ func testGetBloombitsProofs(t *testing.T, protocol int) { trie.Prove(key, 0, &proofs.Proofs) // Send the proof request and verify the response - cost := server.peer.peer.GetRequestCost(GetHelperTrieProofsMsg, len(requests)) + cost := server.peer.speer.getRequestCost(GetHelperTrieProofsMsg, len(requests)) sendRequest(server.peer.app, GetHelperTrieProofsMsg, 42, requests) if err := expectResponse(server.peer.app, protocol, HelperTrieProofsMsg, 42, testBufLimit, cost, proofs); err != nil { t.Errorf("bit %d: proofs mismatch: %v", bit, err) @@ -534,10 +535,10 @@ func testTransactionStatus(t *testing.T, protocol int) { reqID++ var cost uint64 if send { - cost = server.peer.peer.GetRequestCost(SendTxV2Msg, 1) + cost = server.peer.speer.getRequestCost(SendTxV2Msg, 1) sendRequest(server.peer.app, SendTxV2Msg, reqID, types.Transactions{tx}) } else { - cost = server.peer.peer.GetRequestCost(GetTxStatusMsg, 1) + cost = server.peer.speer.getRequestCost(GetTxStatusMsg, 1) sendRequest(server.peer.app, GetTxStatusMsg, reqID, []common.Hash{tx.Hash()}) } if err := expectResponse(server.peer.app, protocol, TxStatusMsg, reqID, testBufLimit, cost, []light.TxStatus{expStatus}); err != nil { diff --git a/les/peer.go b/les/peer.go index 41ea7194bc..a246067e16 100644 --- a/les/peer.go +++ b/les/peer.go @@ -857,40 +857,6 @@ func (p *clientPeer) updateCapacity(cap uint64) { } -// freezeClient temporarily puts the client in a frozen state which means all -// unprocessed and subsequent requests are dropped. Unfreezing happens automatically -// after a short time if the client's buffer value is at least in the slightly positive -// region. The client is also notified about being frozen/unfrozen with a Stop/Resume -// message. -func (p *clientPeer) freezeClient() { - if p.version < lpv3 { - // if Stop/Resume is not supported then just drop the peer after setting - // its frozen status permanently - atomic.StoreUint32(&p.frozen, 1) - p.Peer.Disconnect(p2p.DiscUselessPeer) - return - } - if atomic.SwapUint32(&p.frozen, 1) == 0 { - go func() { - p.sendStop() - time.Sleep(freezeTimeBase + time.Duration(rand.Int63n(int64(freezeTimeRandom)))) - for { - bufValue, bufLimit := p.fcClient.BufferStatus() - if bufLimit == 0 { - return - } - if bufValue <= bufLimit/8 { - time.Sleep(freezeCheckPeriod) - } else { - atomic.StoreUint32(&p.frozen, 0) - p.sendResume(bufValue) - break - } - } - }() - } -} - // Handshake executes the les protocol handshake, negotiating version number, // network IDs, difficulties, head and genesis blocks. func (p *clientPeer) Handshake(td *big.Int, head common.Hash, headNum uint64, genesis common.Hash, server *LesServer) error { @@ -1016,7 +982,7 @@ func (ps *clientPeerSet) unSubscribe(sub clientPeerSubscriber) { // register adds a new peer into the peer set, or returns an error if the // peer is already known. -func (ps *clientPeerSet) register(peer *clientPeer) error { +func (ps *clientPeerSet) register(p *clientPeer) error { ps.lock.Lock() if ps.closed { ps.lock.Unlock() @@ -1042,7 +1008,7 @@ func (ps *clientPeerSet) register(peer *clientPeer) error { // unregister removes a remote peer from the peer set, disabling any further // actions to/from that particular entity. It also initiates disconnection // at the networking layer. -func (ps *clientPeerSet) unregister(id string) error { +func (ps *clientPeerSet) unregister(p *clientPeer) error { ps.lock.Lock() if _, ok := ps.active[p.id]; !ok { ps.lock.Unlock() @@ -1125,7 +1091,7 @@ func (ps *clientPeerSet) allPeers() []*clientPeer { ps.lock.RLock() defer ps.lock.RUnlock() - list := make([]*clientPeer, 0, len(ps.peers)) + list := make([]*clientPeer, 0, len(ps.active)) for _, p := range ps.active { list = append(list, p) } @@ -1197,7 +1163,7 @@ func (ps *serverPeerSet) unSubscribe(sub serverPeerSubscriber) { // register adds a new server peer into the set, or returns an error if the // peer is already known. -func (ps *serverPeerSet) register(peer *serverPeer) error { +func (ps *serverPeerSet) register(p *serverPeer) error { ps.lock.Lock() if ps.closed { ps.lock.Unlock() @@ -1223,7 +1189,7 @@ func (ps *serverPeerSet) register(peer *serverPeer) error { // unregister removes a remote peer from the active set, disabling any further // actions to/from that particular entity. It also initiates disconnection at // the networking layer. -func (ps *serverPeerSet) unregister(id string) error { +func (ps *serverPeerSet) unregister(p *serverPeer) error { ps.lock.Lock() if _, ok := ps.active[p.id]; !ok { ps.lock.Unlock() @@ -1326,7 +1292,7 @@ func (ps *serverPeerSet) allPeers() []*serverPeer { defer ps.lock.RUnlock() list := make([]*serverPeer, 0, len(ps.active)) - for _, p := range ps.peers { + for _, p := range ps.active { list = append(list, p) } return list diff --git a/les/peer_test.go b/les/peer_test.go index 59a2ad7009..a5da9a7a96 100644 --- a/les/peer_test.go +++ b/les/peer_test.go @@ -85,7 +85,7 @@ func TestPeerSubscription(t *testing.T) { checkIds([]string{peer.id}) checkPeers(sub.regCh) - peers.unregister(peer.id) + peers.unregister(peer) checkIds([]string{}) checkPeers(sub.unregCh) } diff --git a/les/server.go b/les/server.go index d6225a9c05..2430d1914a 100644 --- a/les/server.go +++ b/les/server.go @@ -158,7 +158,7 @@ func (s *LesServer) APIs() []rpc.API { { Namespace: "lespay", Version: "1.0", - Service: NewPrivateLespayAPI(s.lesCommons.peers, nil, s.srvr.DiscV5, s.tokenSale), + Service: NewPrivateLespayAPI(s.peers, nil, nil, s.srvr.DiscV5, s.tokenSale), Public: false, }, } diff --git a/les/server_handler.go b/les/server_handler.go index 07bf821102..b928c3fb0b 100644 --- a/les/server_handler.go +++ b/les/server_handler.go @@ -142,12 +142,12 @@ func (h *serverHandler) handle(p *clientPeer) error { ) p.activate = func() { // Register the peer locally - if err := h.server.peers.Register(p); err != nil { + if err := h.server.peers.register(p); err != nil { h.server.clientPool.disconnect(p) p.Log().Error("Light Ethereum peer registration failed", "err", err) return } - clientConnectionGauge.Update(int64(h.server.peers.Len())) + clientConnectionGauge.Update(int64(h.server.peers.len())) connectedAt = mclock.Now() wg = new(sync.WaitGroup) p.active = true @@ -157,7 +157,7 @@ func (h *serverHandler) handle(p *clientPeer) error { if p.version < lpv4 { h.server.peers.disconnect(p.id) } - clientConnectionGauge.Update(int64(h.server.peers.Len())) + clientConnectionGauge.Update(int64(h.server.peers.len())) connectionTimer.Update(time.Duration(mclock.Now() - connectedAt)) p.active = false } @@ -171,7 +171,7 @@ func (h *serverHandler) handle(p *clientPeer) error { return err } else if capacity != p.fcParams.MinRecharge { if p.version < lpv4 { - h.server.peers.Disconnect(p.id) + h.server.peers.disconnect(p.id) } else { p.updateCapacity(capacity) } @@ -417,8 +417,6 @@ func (h *serverHandler) handleMsg(p *clientPeer, wg *sync.WaitGroup) error { first = false } reply := p.replyBlockHeaders(req.ReqID, headers) - sendResponse(req.ReqID, query.Amount, p.replyBlockHeaders(req.ReqID, headers), task.done()) - reply := p.ReplyBlockHeaders(req.ReqID, headers) sendResponse(req.ReqID, query.Amount, reply, task.done()) if metrics.EnabledExpensive { miscOutHeaderPacketsMeter.Mark(1) @@ -894,7 +892,7 @@ func (h *serverHandler) handleMsg(p *clientPeer, wg *sync.WaitGroup) error { miscOutLespayTrafficMeter.Mark(int64(len(reply))) } p.queueSend(func() { - p.ReplyLespay(req.ReqID, reply, delay) + p.replyLespay(req.ReqID, reply, delay) }) }, }) { diff --git a/les/test_helper.go b/les/test_helper.go index b53c28f2e5..95e91371ca 100644 --- a/les/test_helper.go +++ b/les/test_helper.go @@ -309,7 +309,8 @@ func newTestPeer(t *testing.T, name string, version int, handler *serverHandler, // 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) + cpeer := newClientPeer(version, NetworkId, p2p.NewPeer(id, name, nil), net) + speer := newServerPeer(version, NetworkId, false, p2p.NewPeer(id, name, nil), app) // Start the peer on a new thread errCh := make(chan error, 1) @@ -317,13 +318,14 @@ func newTestPeer(t *testing.T, name string, version int, handler *serverHandler, select { case <-handler.closeCh: errCh <- p2p.DiscQuitting - case errCh <- handler.handle(peer): + case errCh <- handler.handle(cpeer): } }() tp := &testPeer{ app: app, net: net, - cpeer: peer, + cpeer: cpeer, + speer: speer, } // Execute any implicitly requested handshakes and return if shake { @@ -395,7 +397,7 @@ func (p *testPeer) handshake(t *testing.T, td *big.Int, head common.Hash, headNu expList = expList.add("serveStateSince", uint64(0)) expList = expList.add("serveRecentState", uint64(core.TriesInMemory-4)) expList = expList.add("txRelay", nil) - if p.peer.version >= lpv4 { + if p.cpeer.version >= lpv4 { expList = expList.add("flowControl/BL", uint64(0)) expList = expList.add("flowControl/MRR", uint64(0)) } else { @@ -413,10 +415,11 @@ func (p *testPeer) handshake(t *testing.T, td *big.Int, head common.Hash, headNu p.cpeer.fcParams = flowcontrol.ServerParams{ BufLimit: testBufLimit, MinRecharge: testBufRecharge, + } } func (p *testPeer) expectCapUpdate(t *testing.T) { - if p.peer.version >= lpv4 { + if p.cpeer.version >= lpv4 { var expList keyValueList expList = expList.add("flowControl/BL", testBufLimit) expList = expList.add("flowControl/MRR", testBufRecharge)