mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-20 10:52:25 +00:00
les: fixed compile errors
This commit is contained in:
parent
cf796e5773
commit
9be4810e77
11 changed files with 61 additions and 84 deletions
37
les/api.go
37
les/api.go
|
|
@ -296,7 +296,7 @@ func (api *PrivateDebugAPI) FreezeClient(id enode.ID) error {
|
||||||
if c == nil {
|
if c == nil {
|
||||||
return fmt.Errorf("client %064x is not connected", id[:])
|
return fmt.Errorf("client %064x is not connected", id[:])
|
||||||
}
|
}
|
||||||
c.peer.freezeClient()
|
c.peer.freeze()
|
||||||
return nil
|
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
|
// PrivateLespayAPI provides an API to use the LESpay commands of either the local or a remote server
|
||||||
type PrivateLespayAPI struct {
|
type PrivateLespayAPI struct {
|
||||||
peerSet *peerSet
|
clientPeerSet *clientPeerSet
|
||||||
|
serverPeerSet *serverPeerSet
|
||||||
clientHandler *clientHandler
|
clientHandler *clientHandler
|
||||||
dht *discv5.Network
|
dht *discv5.Network
|
||||||
tokenSale *tokenSale
|
tokenSale *tokenSale
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewPrivateLespayAPI creates a new LESPAY API.
|
// 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{
|
return &PrivateLespayAPI{
|
||||||
peerSet: peerSet,
|
clientPeerSet: clientPeerSet,
|
||||||
|
serverPeerSet: serverPeerSet,
|
||||||
clientHandler: clientHandler,
|
clientHandler: clientHandler,
|
||||||
dht: dht,
|
dht: dht,
|
||||||
tokenSale: tokenSale,
|
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.
|
// 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) {
|
func (api *PrivateLespayAPI) makeCall(ctx context.Context, remote bool, nodeStr string, cmd []byte) ([]byte, error) {
|
||||||
var (
|
var (
|
||||||
id enode.ID
|
id enode.ID
|
||||||
freeID string
|
freeID string
|
||||||
peer *peer
|
clientPeer *clientPeer
|
||||||
node *enode.Node
|
serverPeer *serverPeer
|
||||||
err error
|
node *enode.Node
|
||||||
|
err error
|
||||||
)
|
)
|
||||||
if nodeStr != "" {
|
if nodeStr != "" {
|
||||||
if id, err = enode.ParseID(nodeStr); err == nil {
|
if id, err = enode.ParseID(nodeStr); err == nil {
|
||||||
if peer = api.peerSet.Peer(peerIdToString(id)); peer == nil {
|
if api.clientPeerSet != nil {
|
||||||
return nil, errors.New("peer not connected")
|
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 {
|
} else {
|
||||||
var err error
|
var err error
|
||||||
if node, err = enode.Parse(enode.ValidSchemes, nodeStr); err == nil {
|
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
|
cancelFn func() bool
|
||||||
)
|
)
|
||||||
delivered := make(chan struct{})
|
delivered := make(chan struct{})
|
||||||
if peer != nil {
|
if serverPeer != nil {
|
||||||
// remote call to a connected peer through LES
|
// remote call to a connected peer through LES
|
||||||
if api.clientHandler == nil {
|
if api.clientHandler == nil {
|
||||||
return nil, errors.New("client handler not available")
|
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
|
reply = r
|
||||||
close(delivered)
|
close(delivered)
|
||||||
return reply != nil
|
return reply != nil
|
||||||
|
|
|
||||||
|
|
@ -212,7 +212,7 @@ func (s *LightEthereum) APIs() []rpc.API {
|
||||||
{
|
{
|
||||||
Namespace: "lespay",
|
Namespace: "lespay",
|
||||||
Version: "1.0",
|
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,
|
Public: false,
|
||||||
},
|
},
|
||||||
}...)
|
}...)
|
||||||
|
|
|
||||||
|
|
@ -391,12 +391,12 @@ func (h *clientHandler) handleMsg(p *serverPeer) error {
|
||||||
// makeLespayCall sends a lespay command through an LES connection and registers
|
// makeLespayCall sends a lespay command through an LES connection and registers
|
||||||
// a response handler. It returns a cancel function that removes the response
|
// 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.
|
// 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()
|
reqID := genReqID()
|
||||||
h.lespayReplyLock.Lock()
|
h.lespayReplyLock.Lock()
|
||||||
h.lespayReplyHandlers[reqID] = handler
|
h.lespayReplyHandlers[reqID] = handler
|
||||||
h.lespayReplyLock.Unlock()
|
h.lespayReplyLock.Unlock()
|
||||||
if p.SendLespay(reqID, cmd) != nil {
|
if p.sendLespay(reqID, cmd) != nil {
|
||||||
h.lespayReplyLock.Lock()
|
h.lespayReplyLock.Lock()
|
||||||
delete(h.lespayReplyHandlers, reqID)
|
delete(h.lespayReplyHandlers, reqID)
|
||||||
h.lespayReplyLock.Unlock()
|
h.lespayReplyLock.Unlock()
|
||||||
|
|
|
||||||
|
|
@ -123,7 +123,7 @@ type clientPoolPeer interface {
|
||||||
ID() enode.ID
|
ID() enode.ID
|
||||||
freeClientId() string
|
freeClientId() string
|
||||||
updateCapacity(uint64)
|
updateCapacity(uint64)
|
||||||
freezeClient()
|
freeze()
|
||||||
}
|
}
|
||||||
|
|
||||||
// clientInfo represents a connected client
|
// 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
|
// 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.
|
// Short circuit if client pool is already closed.
|
||||||
if f.closed {
|
if f.closed {
|
||||||
return
|
return
|
||||||
|
|
|
||||||
|
|
@ -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) {
|
func testClientPool(t *testing.T, activeLimit, clientCount, paidCount int, randomDisconnect bool) {
|
||||||
rand.Seed(time.Now().UnixNano())
|
rand.Seed(time.Now().UnixNano())
|
||||||
|
|
|
||||||
|
|
@ -177,7 +177,7 @@ func testGetBlockHeaders(t *testing.T, protocol int) {
|
||||||
// Send the hash request and verify the response
|
// Send the hash request and verify the response
|
||||||
reqID++
|
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)
|
sendRequest(server.peer.app, GetBlockHeadersMsg, reqID, tt.query)
|
||||||
if err := expectResponse(server.peer.app, protocol, BlockHeadersMsg, reqID, testBufLimit, cost, headers); err != nil {
|
if err := expectResponse(server.peer.app, protocol, BlockHeadersMsg, reqID, testBufLimit, cost, headers); err != nil {
|
||||||
t.Errorf("test %d: headers mismatch: %v", i, err)
|
t.Errorf("test %d: headers mismatch: %v", i, err)
|
||||||
|
|
@ -255,7 +255,7 @@ func testGetBlockBodies(t *testing.T, protocol int) {
|
||||||
reqID++
|
reqID++
|
||||||
|
|
||||||
// Send the hash request and verify the response
|
// 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)
|
sendRequest(server.peer.app, GetBlockBodiesMsg, reqID, hashes)
|
||||||
if err := expectResponse(server.peer.app, protocol, BlockBodiesMsg, reqID, testBufLimit, cost, bodies); err != nil {
|
if err := expectResponse(server.peer.app, protocol, BlockBodiesMsg, reqID, testBufLimit, cost, bodies); err != nil {
|
||||||
t.Errorf("test %d: bodies mismatch: %v", i, err)
|
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)
|
sendRequest(server.peer.app, GetCodeMsg, 42, codereqs)
|
||||||
if err := expectResponse(server.peer.app, protocol, CodeMsg, 42, testBufLimit, cost, codes); err != nil {
|
if err := expectResponse(server.peer.app, protocol, CodeMsg, 42, testBufLimit, cost, codes); err != nil {
|
||||||
t.Errorf("codes mismatch: %v", err)
|
t.Errorf("codes mismatch: %v", err)
|
||||||
|
|
@ -308,7 +308,7 @@ func testGetStaleCode(t *testing.T, protocol int) {
|
||||||
BHash: bc.GetHeaderByNumber(number).Hash(),
|
BHash: bc.GetHeaderByNumber(number).Hash(),
|
||||||
AccKey: crypto.Keccak256(testContractAddr[:]),
|
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})
|
sendRequest(server.peer.app, GetCodeMsg, 42, []*CodeReq{req})
|
||||||
if err := expectResponse(server.peer.app, protocol, CodeMsg, 42, testBufLimit, cost, expected); err != nil {
|
if err := expectResponse(server.peer.app, protocol, CodeMsg, 42, testBufLimit, cost, expected); err != nil {
|
||||||
t.Errorf("codes mismatch: %v", err)
|
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()))
|
receipts = append(receipts, rawdb.ReadRawReceipts(server.db, block.Hash(), block.NumberU64()))
|
||||||
}
|
}
|
||||||
// Send the hash request and verify the response
|
// 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)
|
sendRequest(server.peer.app, GetReceiptsMsg, 42, hashes)
|
||||||
if err := expectResponse(server.peer.app, protocol, ReceiptsMsg, 42, testBufLimit, cost, receipts); err != nil {
|
if err := expectResponse(server.peer.app, protocol, ReceiptsMsg, 42, testBufLimit, cost, receipts); err != nil {
|
||||||
t.Errorf("receipts mismatch: %v", err)
|
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
|
// 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)
|
sendRequest(server.peer.app, GetProofsV2Msg, 42, proofreqs)
|
||||||
if err := expectResponse(server.peer.app, protocol, ProofsV2Msg, 42, testBufLimit, cost, proofsV2.NodeList()); err != nil {
|
if err := expectResponse(server.peer.app, protocol, ProofsV2Msg, 42, testBufLimit, cost, proofsV2.NodeList()); err != nil {
|
||||||
t.Errorf("proofs mismatch: %v", err)
|
t.Errorf("proofs mismatch: %v", err)
|
||||||
|
|
@ -410,6 +410,7 @@ func testGetStaleProof(t *testing.T, protocol int) {
|
||||||
t.Prove(account, 0, proofsV2)
|
t.Prove(account, 0, proofsV2)
|
||||||
expected = proofsV2.NodeList()
|
expected = proofsV2.NodeList()
|
||||||
}
|
}
|
||||||
|
cost := server.peer.speer.getRequestCost(GetProofsV2Msg, 1)
|
||||||
if err := expectResponse(server.peer.app, protocol, ProofsV2Msg, 42, testBufLimit, cost, expected); err != nil {
|
if err := expectResponse(server.peer.app, protocol, ProofsV2Msg, 42, testBufLimit, cost, expected); err != nil {
|
||||||
t.Errorf("codes mismatch: %v", err)
|
t.Errorf("codes mismatch: %v", err)
|
||||||
}
|
}
|
||||||
|
|
@ -461,7 +462,7 @@ func testGetCHTProofs(t *testing.T, protocol int) {
|
||||||
AuxReq: auxHeader,
|
AuxReq: auxHeader,
|
||||||
}}
|
}}
|
||||||
// Send the proof request and verify the response
|
// 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)
|
sendRequest(server.peer.app, GetHelperTrieProofsMsg, 42, requestsV2)
|
||||||
if err := expectResponse(server.peer.app, protocol, HelperTrieProofsMsg, 42, testBufLimit, cost, proofsV2); err != nil {
|
if err := expectResponse(server.peer.app, protocol, HelperTrieProofsMsg, 42, testBufLimit, cost, proofsV2); err != nil {
|
||||||
t.Errorf("proofs mismatch: %v", err)
|
t.Errorf("proofs mismatch: %v", err)
|
||||||
|
|
@ -510,7 +511,7 @@ func testGetBloombitsProofs(t *testing.T, protocol int) {
|
||||||
trie.Prove(key, 0, &proofs.Proofs)
|
trie.Prove(key, 0, &proofs.Proofs)
|
||||||
|
|
||||||
// Send the proof request and verify the response
|
// 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)
|
sendRequest(server.peer.app, GetHelperTrieProofsMsg, 42, requests)
|
||||||
if err := expectResponse(server.peer.app, protocol, HelperTrieProofsMsg, 42, testBufLimit, cost, proofs); err != nil {
|
if err := expectResponse(server.peer.app, protocol, HelperTrieProofsMsg, 42, testBufLimit, cost, proofs); err != nil {
|
||||||
t.Errorf("bit %d: proofs mismatch: %v", bit, err)
|
t.Errorf("bit %d: proofs mismatch: %v", bit, err)
|
||||||
|
|
@ -534,10 +535,10 @@ func testTransactionStatus(t *testing.T, protocol int) {
|
||||||
reqID++
|
reqID++
|
||||||
var cost uint64
|
var cost uint64
|
||||||
if send {
|
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})
|
sendRequest(server.peer.app, SendTxV2Msg, reqID, types.Transactions{tx})
|
||||||
} else {
|
} 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()})
|
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 {
|
if err := expectResponse(server.peer.app, protocol, TxStatusMsg, reqID, testBufLimit, cost, []light.TxStatus{expStatus}); err != nil {
|
||||||
|
|
|
||||||
46
les/peer.go
46
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,
|
// Handshake executes the les protocol handshake, negotiating version number,
|
||||||
// network IDs, difficulties, head and genesis blocks.
|
// 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 {
|
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
|
// register adds a new peer into the peer set, or returns an error if the
|
||||||
// peer is already known.
|
// peer is already known.
|
||||||
func (ps *clientPeerSet) register(peer *clientPeer) error {
|
func (ps *clientPeerSet) register(p *clientPeer) error {
|
||||||
ps.lock.Lock()
|
ps.lock.Lock()
|
||||||
if ps.closed {
|
if ps.closed {
|
||||||
ps.lock.Unlock()
|
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
|
// unregister removes a remote peer from the peer set, disabling any further
|
||||||
// actions to/from that particular entity. It also initiates disconnection
|
// actions to/from that particular entity. It also initiates disconnection
|
||||||
// at the networking layer.
|
// at the networking layer.
|
||||||
func (ps *clientPeerSet) unregister(id string) error {
|
func (ps *clientPeerSet) unregister(p *clientPeer) error {
|
||||||
ps.lock.Lock()
|
ps.lock.Lock()
|
||||||
if _, ok := ps.active[p.id]; !ok {
|
if _, ok := ps.active[p.id]; !ok {
|
||||||
ps.lock.Unlock()
|
ps.lock.Unlock()
|
||||||
|
|
@ -1125,7 +1091,7 @@ func (ps *clientPeerSet) allPeers() []*clientPeer {
|
||||||
ps.lock.RLock()
|
ps.lock.RLock()
|
||||||
defer ps.lock.RUnlock()
|
defer ps.lock.RUnlock()
|
||||||
|
|
||||||
list := make([]*clientPeer, 0, len(ps.peers))
|
list := make([]*clientPeer, 0, len(ps.active))
|
||||||
for _, p := range ps.active {
|
for _, p := range ps.active {
|
||||||
list = append(list, p)
|
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
|
// register adds a new server peer into the set, or returns an error if the
|
||||||
// peer is already known.
|
// peer is already known.
|
||||||
func (ps *serverPeerSet) register(peer *serverPeer) error {
|
func (ps *serverPeerSet) register(p *serverPeer) error {
|
||||||
ps.lock.Lock()
|
ps.lock.Lock()
|
||||||
if ps.closed {
|
if ps.closed {
|
||||||
ps.lock.Unlock()
|
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
|
// unregister removes a remote peer from the active set, disabling any further
|
||||||
// actions to/from that particular entity. It also initiates disconnection at
|
// actions to/from that particular entity. It also initiates disconnection at
|
||||||
// the networking layer.
|
// the networking layer.
|
||||||
func (ps *serverPeerSet) unregister(id string) error {
|
func (ps *serverPeerSet) unregister(p *serverPeer) error {
|
||||||
ps.lock.Lock()
|
ps.lock.Lock()
|
||||||
if _, ok := ps.active[p.id]; !ok {
|
if _, ok := ps.active[p.id]; !ok {
|
||||||
ps.lock.Unlock()
|
ps.lock.Unlock()
|
||||||
|
|
@ -1326,7 +1292,7 @@ func (ps *serverPeerSet) allPeers() []*serverPeer {
|
||||||
defer ps.lock.RUnlock()
|
defer ps.lock.RUnlock()
|
||||||
|
|
||||||
list := make([]*serverPeer, 0, len(ps.active))
|
list := make([]*serverPeer, 0, len(ps.active))
|
||||||
for _, p := range ps.peers {
|
for _, p := range ps.active {
|
||||||
list = append(list, p)
|
list = append(list, p)
|
||||||
}
|
}
|
||||||
return list
|
return list
|
||||||
|
|
|
||||||
|
|
@ -85,7 +85,7 @@ func TestPeerSubscription(t *testing.T) {
|
||||||
checkIds([]string{peer.id})
|
checkIds([]string{peer.id})
|
||||||
checkPeers(sub.regCh)
|
checkPeers(sub.regCh)
|
||||||
|
|
||||||
peers.unregister(peer.id)
|
peers.unregister(peer)
|
||||||
checkIds([]string{})
|
checkIds([]string{})
|
||||||
checkPeers(sub.unregCh)
|
checkPeers(sub.unregCh)
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -158,7 +158,7 @@ func (s *LesServer) APIs() []rpc.API {
|
||||||
{
|
{
|
||||||
Namespace: "lespay",
|
Namespace: "lespay",
|
||||||
Version: "1.0",
|
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,
|
Public: false,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -142,12 +142,12 @@ func (h *serverHandler) handle(p *clientPeer) error {
|
||||||
)
|
)
|
||||||
p.activate = func() {
|
p.activate = func() {
|
||||||
// Register the peer locally
|
// 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)
|
h.server.clientPool.disconnect(p)
|
||||||
p.Log().Error("Light Ethereum peer registration failed", "err", err)
|
p.Log().Error("Light Ethereum peer registration failed", "err", err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
clientConnectionGauge.Update(int64(h.server.peers.Len()))
|
clientConnectionGauge.Update(int64(h.server.peers.len()))
|
||||||
connectedAt = mclock.Now()
|
connectedAt = mclock.Now()
|
||||||
wg = new(sync.WaitGroup)
|
wg = new(sync.WaitGroup)
|
||||||
p.active = true
|
p.active = true
|
||||||
|
|
@ -157,7 +157,7 @@ func (h *serverHandler) handle(p *clientPeer) error {
|
||||||
if p.version < lpv4 {
|
if p.version < lpv4 {
|
||||||
h.server.peers.disconnect(p.id)
|
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))
|
connectionTimer.Update(time.Duration(mclock.Now() - connectedAt))
|
||||||
p.active = false
|
p.active = false
|
||||||
}
|
}
|
||||||
|
|
@ -171,7 +171,7 @@ func (h *serverHandler) handle(p *clientPeer) error {
|
||||||
return err
|
return err
|
||||||
} else if capacity != p.fcParams.MinRecharge {
|
} else if capacity != p.fcParams.MinRecharge {
|
||||||
if p.version < lpv4 {
|
if p.version < lpv4 {
|
||||||
h.server.peers.Disconnect(p.id)
|
h.server.peers.disconnect(p.id)
|
||||||
} else {
|
} else {
|
||||||
p.updateCapacity(capacity)
|
p.updateCapacity(capacity)
|
||||||
}
|
}
|
||||||
|
|
@ -417,8 +417,6 @@ func (h *serverHandler) handleMsg(p *clientPeer, wg *sync.WaitGroup) error {
|
||||||
first = false
|
first = false
|
||||||
}
|
}
|
||||||
reply := p.replyBlockHeaders(req.ReqID, headers)
|
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())
|
sendResponse(req.ReqID, query.Amount, reply, task.done())
|
||||||
if metrics.EnabledExpensive {
|
if metrics.EnabledExpensive {
|
||||||
miscOutHeaderPacketsMeter.Mark(1)
|
miscOutHeaderPacketsMeter.Mark(1)
|
||||||
|
|
@ -894,7 +892,7 @@ func (h *serverHandler) handleMsg(p *clientPeer, wg *sync.WaitGroup) error {
|
||||||
miscOutLespayTrafficMeter.Mark(int64(len(reply)))
|
miscOutLespayTrafficMeter.Mark(int64(len(reply)))
|
||||||
}
|
}
|
||||||
p.queueSend(func() {
|
p.queueSend(func() {
|
||||||
p.ReplyLespay(req.ReqID, reply, delay)
|
p.replyLespay(req.ReqID, reply, delay)
|
||||||
})
|
})
|
||||||
},
|
},
|
||||||
}) {
|
}) {
|
||||||
|
|
|
||||||
|
|
@ -309,7 +309,8 @@ func newTestPeer(t *testing.T, name string, version int, handler *serverHandler,
|
||||||
// Generate a random id and create the peer
|
// Generate a random id and create the peer
|
||||||
var id enode.ID
|
var id enode.ID
|
||||||
rand.Read(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
|
// Start the peer on a new thread
|
||||||
errCh := make(chan error, 1)
|
errCh := make(chan error, 1)
|
||||||
|
|
@ -317,13 +318,14 @@ func newTestPeer(t *testing.T, name string, version int, handler *serverHandler,
|
||||||
select {
|
select {
|
||||||
case <-handler.closeCh:
|
case <-handler.closeCh:
|
||||||
errCh <- p2p.DiscQuitting
|
errCh <- p2p.DiscQuitting
|
||||||
case errCh <- handler.handle(peer):
|
case errCh <- handler.handle(cpeer):
|
||||||
}
|
}
|
||||||
}()
|
}()
|
||||||
tp := &testPeer{
|
tp := &testPeer{
|
||||||
app: app,
|
app: app,
|
||||||
net: net,
|
net: net,
|
||||||
cpeer: peer,
|
cpeer: cpeer,
|
||||||
|
speer: speer,
|
||||||
}
|
}
|
||||||
// Execute any implicitly requested handshakes and return
|
// Execute any implicitly requested handshakes and return
|
||||||
if shake {
|
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("serveStateSince", uint64(0))
|
||||||
expList = expList.add("serveRecentState", uint64(core.TriesInMemory-4))
|
expList = expList.add("serveRecentState", uint64(core.TriesInMemory-4))
|
||||||
expList = expList.add("txRelay", nil)
|
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/BL", uint64(0))
|
||||||
expList = expList.add("flowControl/MRR", uint64(0))
|
expList = expList.add("flowControl/MRR", uint64(0))
|
||||||
} else {
|
} else {
|
||||||
|
|
@ -413,10 +415,11 @@ func (p *testPeer) handshake(t *testing.T, td *big.Int, head common.Hash, headNu
|
||||||
p.cpeer.fcParams = flowcontrol.ServerParams{
|
p.cpeer.fcParams = flowcontrol.ServerParams{
|
||||||
BufLimit: testBufLimit,
|
BufLimit: testBufLimit,
|
||||||
MinRecharge: testBufRecharge,
|
MinRecharge: testBufRecharge,
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func (p *testPeer) expectCapUpdate(t *testing.T) {
|
func (p *testPeer) expectCapUpdate(t *testing.T) {
|
||||||
if p.peer.version >= lpv4 {
|
if p.cpeer.version >= lpv4 {
|
||||||
var expList keyValueList
|
var expList keyValueList
|
||||||
expList = expList.add("flowControl/BL", testBufLimit)
|
expList = expList.add("flowControl/BL", testBufLimit)
|
||||||
expList = expList.add("flowControl/MRR", testBufRecharge)
|
expList = expList.add("flowControl/MRR", testBufRecharge)
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue