les: lespay added to les/4

This commit is contained in:
Zsolt Felfoldi 2019-11-24 20:30:12 +01:00
parent 1ddc30da95
commit 6ac3b2b81b
5 changed files with 50 additions and 6 deletions

View file

@ -40,6 +40,8 @@ var (
miscInTxsTrafficMeter = metrics.NewRegisteredMeter("les/misc/in/traffic/txs", nil)
miscInTxStatusPacketsMeter = metrics.NewRegisteredMeter("les/misc/in/packets/txStatus", nil)
miscInTxStatusTrafficMeter = metrics.NewRegisteredMeter("les/misc/in/traffic/txStatus", nil)
miscInLespayPacketsMeter = metrics.NewRegisteredMeter("les/misc/in/packets/lespay", nil)
miscInLespayTrafficMeter = metrics.NewRegisteredMeter("les/misc/in/traffic/lespay", nil)
miscOutPacketsMeter = metrics.NewRegisteredMeter("les/misc/out/packets/total", nil)
miscOutTrafficMeter = metrics.NewRegisteredMeter("les/misc/out/traffic/total", nil)
@ -59,6 +61,8 @@ var (
miscOutTxsTrafficMeter = metrics.NewRegisteredMeter("les/misc/out/traffic/txs", nil)
miscOutTxStatusPacketsMeter = metrics.NewRegisteredMeter("les/misc/out/packets/txStatus", nil)
miscOutTxStatusTrafficMeter = metrics.NewRegisteredMeter("les/misc/out/traffic/txStatus", nil)
miscOutLespayPacketsMeter = metrics.NewRegisteredMeter("les/misc/out/packets/lespay", nil)
miscOutLespayTrafficMeter = metrics.NewRegisteredMeter("les/misc/out/traffic/lespay", nil)
miscServingTimeHeaderTimer = metrics.NewRegisteredTimer("les/misc/serve/header", nil)
miscServingTimeBodyTimer = metrics.NewRegisteredTimer("les/misc/serve/body", nil)
@ -68,6 +72,7 @@ var (
miscServingTimeHelperTrieTimer = metrics.NewRegisteredTimer("les/misc/serve/helperTrie", nil)
miscServingTimeTxTimer = metrics.NewRegisteredTimer("les/misc/serve/txs", nil)
miscServingTimeTxStatusTimer = metrics.NewRegisteredTimer("les/misc/serve/txStatus", nil)
miscServingTimeLespayTimer = metrics.NewRegisteredTimer("les/misc/serve/lespay", nil)
connectionTimer = metrics.NewRegisteredTimer("les/connection/duration", nil)
serverConnectionGauge = metrics.NewRegisteredGauge("les/connection/server", nil)

View file

@ -501,6 +501,18 @@ func (p *peer) SendTxs(reqID, cost uint64, txs rlp.RawValue) error {
return sendRequest(p.rw, SendTxV2Msg, reqID, cost, txs)
}
// SendLespay sends a set of commands to the service token sale module
func (p *peer) SendLespay(reqID uint64, cmds [][]byte) error {
p.Log().Debug("Sending batch of lespay commands", "size", len(cmds))
return sendRequest(p.rw, LespayMsg, reqID, 0, cmds)
}
// ReplyLespay sends a set of replies to lespay commands
func (p *peer) ReplyLespay(reqID uint64, replies [][]byte) error {
p.Log().Debug("Sending batch of lespay replies", "size", len(replies))
return sendRequest(p.rw, LespayReplyMsg, reqID, 0, replies)
}
type keyValueEntry struct {
Key string
Value rlp.RawValue

View file

@ -33,17 +33,18 @@ import (
const (
lpv2 = 2
lpv3 = 3
lpv4 = 4
)
// Supported versions of the les protocol (first is primary)
var (
ClientProtocolVersions = []uint{lpv2, lpv3}
ServerProtocolVersions = []uint{lpv2, lpv3}
ClientProtocolVersions = []uint{lpv2, lpv3, lpv4}
ServerProtocolVersions = []uint{lpv2, lpv3, lpv4}
AdvertiseProtocolVersions = []uint{lpv2} // clients are searching for the first advertised protocol in the list
)
// Number of implemented message corresponding to different protocol versions.
var ProtocolLengths = map[uint]uint64{lpv2: 22, lpv3: 24}
var ProtocolLengths = map[uint]uint64{lpv2: 22, lpv3: 24, lpv4: 26}
const (
NetworkId = 1
@ -74,6 +75,9 @@ const (
// Protocol messages introduced in LPV3
StopMsg = 0x16
ResumeMsg = 0x17
// Protocol messages introduced in LPV4
LespayMsg = 0x18
LespayReplyMsg = 0x19
)
type requestInfo struct {

View file

@ -118,6 +118,7 @@ func NewLesServer(e *eth.Ethereum, config *eth.Config) (*LesServer, error) {
srv.fcManager.SetCapacityLimits(srv.freeCapacity, srv.maxCapacity, srv.freeCapacity*2)
srv.clientPool = newClientPool(srv.chainDb, srv.minCapacity, srv.freeCapacity, mclock.System{}, func(id enode.ID) { go srv.peers.Unregister(peerIdToString(id)) })
srv.clientPool.setDefaultFactors(priceFactors{0, 1, 1}, priceFactors{0, 1, 1})
srv.tokenSale = newTokenSale(srv.clientPool, 0.1)
checkpoint := srv.latestLocalCheckpoint()
if !checkpoint.Empty() {

View file

@ -824,6 +824,31 @@ func (h *serverHandler) handleMsg(p *peer, wg *sync.WaitGroup) error {
}
}()
}
case LespayMsg:
p.Log().Trace("Received transaction status query request")
if metrics.EnabledExpensive {
miscInLespayPacketsMeter.Mark(1)
miscInLespayTrafficMeter.Mark(int64(msg.Size))
defer func(start time.Time) { miscServingTimeLespayTimer.UpdateSince(start) }(time.Now())
}
var req struct {
ReqID uint64
Cmds [][]byte
}
if err := msg.Decode(&req); err != nil {
clientErrorMeter.Mark(1)
return errResp(ErrDecode, "msg %v: %v", msg, err)
}
replies := h.server.tokenSale.runCommands(req.Cmds, p.ID(), p.freeClientId())
p.ReplyLespay(req.ReqID, replies)
if metrics.EnabledExpensive {
miscOutLespayPacketsMeter.Mark(1)
var size int64
for _, r := range replies {
size += int64(len(r))
}
miscOutLespayTrafficMeter.Mark(size)
}
default:
p.Log().Trace("Received invalid message", "code", msg.Code)
@ -957,9 +982,6 @@ func (h *serverHandler) broadcastHeaders() {
}
func (h *serverHandler) talkRequestHandler(id enode.ID, addr *net.UDPAddr, payload []byte) ([]byte, bool) {
if h.server.tokenSale == nil {
return nil, false
}
var cmds [][]byte
if err := rlp.DecodeBytes(payload, &cmds); err != nil {
return nil, false