From 4b5373b2a5ba2746569255f716a0a66295097796 Mon Sep 17 00:00:00 2001 From: William Banfield Date: Fri, 21 Oct 2022 11:54:22 -0400 Subject: [PATCH] remove wrap from blocksync --- blocksync/msgs.go | 49 ----------------------------- blocksync/reactor.go | 75 ++++++++++---------------------------------- 2 files changed, 16 insertions(+), 108 deletions(-) diff --git a/blocksync/msgs.go b/blocksync/msgs.go index fe1c87d7f..142c38716 100644 --- a/blocksync/msgs.go +++ b/blocksync/msgs.go @@ -19,55 +19,6 @@ const ( BlockResponseMessageFieldKeySize ) -func wrapMsg(pb proto.Message) (proto.Message, error) { - msg := bcproto.Message{} - - switch pb := pb.(type) { - case *bcproto.BlockRequest: - msg.Sum = &bcproto.Message_BlockRequest{BlockRequest: pb} - case *bcproto.BlockResponse: - msg.Sum = &bcproto.Message_BlockResponse{BlockResponse: pb} - case *bcproto.NoBlockResponse: - msg.Sum = &bcproto.Message_NoBlockResponse{NoBlockResponse: pb} - case *bcproto.StatusRequest: - msg.Sum = &bcproto.Message_StatusRequest{StatusRequest: pb} - case *bcproto.StatusResponse: - msg.Sum = &bcproto.Message_StatusResponse{StatusResponse: pb} - default: - return nil, fmt.Errorf("unknown message type %T", pb) - } - - return &msg, nil -} - -// DecodeMsg decodes a Protobuf message. -func DecodeMsg(bz []byte) (proto.Message, error) { - pb := &bcproto.Message{} - - err := proto.Unmarshal(bz, pb) - if err != nil { - return nil, err - } - return UnwrapMessage(pb) -} - -func UnwrapMessage(pb *bcproto.Message) (proto.Message, error) { - switch msg := pb.Sum.(type) { - case *bcproto.Message_BlockRequest: - return msg.BlockRequest, nil - case *bcproto.Message_BlockResponse: - return msg.BlockResponse, nil - case *bcproto.Message_NoBlockResponse: - return msg.NoBlockResponse, nil - case *bcproto.Message_StatusRequest: - return msg.StatusRequest, nil - case *bcproto.Message_StatusResponse: - return msg.StatusResponse, nil - default: - return nil, fmt.Errorf("unknown message type %T", msg) - } -} - // ValidateMsg validates a message. func ValidateMsg(pb proto.Message) error { if pb == nil { diff --git a/blocksync/reactor.go b/blocksync/reactor.go index 9bac4229b..9a1864b38 100644 --- a/blocksync/reactor.go +++ b/blocksync/reactor.go @@ -150,18 +150,12 @@ func (bcR *Reactor) GetChannels() []*p2p.ChannelDescriptor { // AddPeer implements Reactor by sending our state to peer. func (bcR *Reactor) AddPeer(peer p2p.Peer) { - msg, err := wrapMsg(&bcproto.StatusResponse{ - Base: bcR.store.Base(), - Height: bcR.store.Height(), - }) - if err != nil { - bcR.Logger.Error("could not convert msg to protobuf", "err", err) - return - } - peer.Send(p2p.Envelope{ ChannelID: BlocksyncChannel, - Message: msg, + Message: &bcproto.StatusResponse{ + Base: bcR.store.Base(), + Height: bcR.store.Height(), + }, }) // it's OK if send fails. will try later in poolRoutine @@ -187,50 +181,30 @@ func (bcR *Reactor) respondToPeer(msg *bcproto.BlockRequest, return false } - wm, err := wrapMsg(&bcproto.BlockResponse{Block: bl}) - if err != nil { - bcR.Logger.Error("could not convert msg to proto message", "err", err) - return false - } - return src.TrySend(p2p.Envelope{ ChannelID: BlocksyncChannel, - Message: wm, + Message: &bcproto.BlockResponse{Block: bl}, }) } bcR.Logger.Info("Peer asking for a block we don't have", "src", src, "height", msg.Height) - - wm, err := wrapMsg(&bcproto.NoBlockResponse{Height: msg.Height}) - if err != nil { - bcR.Logger.Error("could not convert msg to protobuf", "err", err) - return false - } - return src.TrySend(p2p.Envelope{ ChannelID: BlocksyncChannel, - Message: wm, + Message: &bcproto.NoBlockResponse{Height: msg.Height}, }) } // Receive implements Reactor by handling 4 types of messages (look below). func (bcR *Reactor) Receive(e p2p.Envelope) { - msg, err := UnwrapMessage(e.Message.(*bcproto.Message)) - if err != nil { - bcR.Logger.Error("Error decoding message", "src", e.Src, "chId", e.ChannelID, "err", err) + if err := ValidateMsg(e.Message); err != nil { + bcR.Logger.Error("Peer sent us invalid msg", "peer", e.Src, "msg", e.Message, "err", err) bcR.Switch.StopPeerForError(e.Src, err) return } - if err = ValidateMsg(msg); err != nil { - bcR.Logger.Error("Peer sent us invalid msg", "peer", e.Src, "msg", msg, "err", err) - bcR.Switch.StopPeerForError(e.Src, err) - return - } + bcR.Logger.Debug("Receive", "e.Src", e.Src, "chID", e.ChannelID, "msg", e.Message) - bcR.Logger.Debug("Receive", "e.Src", e.Src, "chID", e.ChannelID, "msg", msg) - - switch msg := msg.(type) { + switch msg := e.Message.(type) { case *bcproto.BlockRequest: bcR.respondToPeer(msg, e.Src) case *bcproto.BlockResponse: @@ -242,18 +216,12 @@ func (bcR *Reactor) Receive(e p2p.Envelope) { bcR.pool.AddBlock(e.Src.ID(), bi, msg.Block.Size()) case *bcproto.StatusRequest: // Send peer our state. - wm, err := wrapMsg(&bcproto.StatusResponse{ - Height: bcR.store.Height(), - Base: bcR.store.Base(), - }) - if err != nil { - bcR.Logger.Error("could not convert msg to proto message", "err", err) - return - } - e.Src.TrySend(p2p.Envelope{ ChannelID: BlocksyncChannel, - Message: wm, + Message: &bcproto.StatusResponse{ + Height: bcR.store.Height(), + Base: bcR.store.Base(), + }, }) case *bcproto.StatusResponse: // Got a peer status. Unverified. @@ -300,14 +268,9 @@ func (bcR *Reactor) poolRoutine(stateSynced bool) { if peer == nil { continue } - wm, err := wrapMsg(&bcproto.BlockRequest{Height: request.Height}) - if err != nil { - bcR.Logger.Error("could not convert msg to proto", "err", err) - continue - } queued := peer.TrySend(p2p.Envelope{ ChannelID: BlocksyncChannel, - Message: wm, + Message: &bcproto.BlockRequest{Height: request.Height}, }) if !queued { bcR.Logger.Debug("Send queue is full, drop block request", "peer", peer.ID(), "height", request.Height) @@ -447,15 +410,9 @@ FOR_LOOP: // BroadcastStatusRequest broadcasts `BlockStore` base and height. func (bcR *Reactor) BroadcastStatusRequest() error { - wm, err := wrapMsg(&bcproto.StatusRequest{}) - if err != nil { - bcR.Logger.Error("could not convert msg to proto message", "err", err) - return fmt.Errorf("could not convert msg to proto message: %w", err) - } bcR.Switch.NewBroadcast(p2p.Envelope{ ChannelID: BlocksyncChannel, - Message: wm, + Message: &bcproto.StatusRequest{}, }) - return nil }