les: use single command format for lespay/reply

This commit is contained in:
Zsolt Felfoldi 2019-11-28 00:45:43 +01:00
parent 8c9f0eb7bf
commit 9c7e7176b0
4 changed files with 19 additions and 25 deletions

View file

@ -407,10 +407,8 @@ func (api *PrivateLespayAPI) makeCall(ctx context.Context, remote bool, nodeStr
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, [][]byte{cmd}, func(replies [][]byte) bool { cancelFn = api.clientHandler.makeLespayCall(peer, cmd, func(r []byte) bool {
if len(replies) == 1 { reply = r
reply = replies[0]
}
close(delivered) close(delivered)
return reply != nil return reply != nil
}) })

View file

@ -41,7 +41,7 @@ type clientHandler struct {
downloader *downloader.Downloader downloader *downloader.Downloader
backend *LightEthereum backend *LightEthereum
lespayReplyHandlers map[uint64]func([][]byte) bool lespayReplyHandlers map[uint64]func([]byte) bool
lespayReplyLock sync.Mutex lespayReplyLock sync.Mutex
closeCh chan struct{} closeCh chan struct{}
@ -54,7 +54,7 @@ func newClientHandler(ulcServers []string, ulcFraction int, checkpoint *params.T
checkpoint: checkpoint, checkpoint: checkpoint,
backend: backend, backend: backend,
closeCh: make(chan struct{}), closeCh: make(chan struct{}),
lespayReplyHandlers: make(map[uint64]func([][]byte) bool), lespayReplyHandlers: make(map[uint64]func([]byte) bool),
} }
if ulcServers != nil { if ulcServers != nil {
ulc, err := newULC(ulcServers, ulcFraction) ulc, err := newULC(ulcServers, ulcFraction)
@ -320,8 +320,8 @@ func (h *clientHandler) handleMsg(p *peer) error {
fmt.Println("LespayReply received") fmt.Println("LespayReply received")
p.Log().Trace("Received tx status response") p.Log().Trace("Received tx status response")
var resp struct { var resp struct {
ReqID uint64 ReqID uint64
Replies [][]byte Reply []byte
} }
if err := msg.Decode(&resp); err != nil { if err := msg.Decode(&resp); err != nil {
fmt.Println("LespayReply decode err", err) fmt.Println("LespayReply decode err", err)
@ -332,7 +332,7 @@ func (h *clientHandler) handleMsg(p *peer) error {
if handler := h.lespayReplyHandlers[resp.ReqID]; handler != nil { if handler := h.lespayReplyHandlers[resp.ReqID]; handler != nil {
fmt.Println("handler found") fmt.Println("handler found")
delete(h.lespayReplyHandlers, resp.ReqID) delete(h.lespayReplyHandlers, resp.ReqID)
responseError = !handler(resp.Replies) responseError = !handler(resp.Reply)
} else { } else {
fmt.Println("handler not found") fmt.Println("handler not found")
responseError = true responseError = true
@ -358,12 +358,12 @@ func (h *clientHandler) handleMsg(p *peer) error {
return nil return nil
} }
func (h *clientHandler) makeLespayCall(p *peer, cmds [][]byte, handler func([][]byte) bool) func() bool { func (h *clientHandler) makeLespayCall(p *peer, cmd []byte, handler func([]byte) 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, cmds) != 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()

View file

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

View file

@ -835,21 +835,17 @@ func (h *serverHandler) handleMsg(p *peer, wg *sync.WaitGroup) error {
} }
var req struct { var req struct {
ReqID uint64 ReqID uint64
Cmds [][]byte Cmd []byte
} }
if err := msg.Decode(&req); err != nil { if err := msg.Decode(&req); err != nil {
clientErrorMeter.Mark(1) clientErrorMeter.Mark(1)
return errResp(ErrDecode, "msg %v: %v", msg, err) return errResp(ErrDecode, "msg %v: %v", msg, err)
} }
replies := h.server.tokenSale.runCommands(req.Cmds, p.ID(), p.freeClientId()) reply := h.server.tokenSale.runCommand(req.Cmd, p.ID(), p.freeClientId())
p.ReplyLespay(req.ReqID, replies) p.ReplyLespay(req.ReqID, reply)
if metrics.EnabledExpensive { if metrics.EnabledExpensive {
miscOutLespayPacketsMeter.Mark(1) miscOutLespayPacketsMeter.Mark(1)
var size int64 miscOutLespayTrafficMeter.Mark(int64(len(reply)))
for _, r := range replies {
size += int64(len(r))
}
miscOutLespayTrafficMeter.Mark(size)
} }
default: default: