Files
node/message/inbound_msg_builder.go

351 lines
7.2 KiB
Go

// Copyright (C) 2019-2025, Lux Industries Inc. All rights reserved.
// See the file LICENSE for licensing terms.
package message
import (
"time"
"github.com/luxfi/ids"
"github.com/luxfi/node/proto/p2p"
"github.com/luxfi/timer/mockable"
)
var _ InboundMsgBuilder = (*inMsgBuilder)(nil)
type InboundMsgBuilder interface {
// Parse reads given bytes as InboundMessage
Parse(
bytes []byte,
nodeID ids.NodeID,
onFinishedHandling func(),
) (InboundMessage, error)
}
type inMsgBuilder struct {
builder *msgBuilder
}
func newInboundBuilder(builder *msgBuilder) InboundMsgBuilder {
return &inMsgBuilder{
builder: builder,
}
}
func (b *inMsgBuilder) Parse(bytes []byte, nodeID ids.NodeID, onFinishedHandling func()) (InboundMessage, error) {
return b.builder.parseInbound(bytes, nodeID, onFinishedHandling)
}
func InboundGetStateSummaryFrontier(
chainID ids.ID,
requestID uint32,
deadline time.Duration,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: GetStateSummaryFrontierOp,
message: &p2p.GetStateSummaryFrontier{
ChainId: chainID[:],
RequestId: requestID,
Deadline: uint64(deadline),
},
expiration: time.Now().Add(deadline),
}
}
func InboundStateSummaryFrontier(
chainID ids.ID,
requestID uint32,
summary []byte,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: StateSummaryFrontierOp,
message: &p2p.StateSummaryFrontier{
ChainId: chainID[:],
RequestId: requestID,
Summary: summary,
},
expiration: mockable.MaxTime,
}
}
func InboundGetAcceptedStateSummary(
chainID ids.ID,
requestID uint32,
heights []uint64,
deadline time.Duration,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: GetAcceptedStateSummaryOp,
message: &p2p.GetAcceptedStateSummary{
ChainId: chainID[:],
RequestId: requestID,
Deadline: uint64(deadline),
Heights: heights,
},
expiration: time.Now().Add(deadline),
}
}
func InboundAcceptedStateSummary(
chainID ids.ID,
requestID uint32,
summaryIDs []ids.ID,
nodeID ids.NodeID,
) InboundMessage {
summaryIDBytes := make([][]byte, len(summaryIDs))
encodeIDs(summaryIDs, summaryIDBytes)
return &inboundMessage{
nodeID: nodeID,
op: AcceptedStateSummaryOp,
message: &p2p.AcceptedStateSummary{
ChainId: chainID[:],
RequestId: requestID,
SummaryIds: summaryIDBytes,
},
expiration: mockable.MaxTime,
}
}
func InboundGetAcceptedFrontier(
chainID ids.ID,
requestID uint32,
deadline time.Duration,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: GetAcceptedFrontierOp,
message: &p2p.GetAcceptedFrontier{
ChainId: chainID[:],
RequestId: requestID,
Deadline: uint64(deadline),
},
expiration: time.Now().Add(deadline),
}
}
func InboundAcceptedFrontier(
chainID ids.ID,
requestID uint32,
containerID ids.ID,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: AcceptedFrontierOp,
message: &p2p.AcceptedFrontier{
ChainId: chainID[:],
RequestId: requestID,
ContainerId: containerID[:],
},
expiration: mockable.MaxTime,
}
}
func InboundGetAccepted(
chainID ids.ID,
requestID uint32,
deadline time.Duration,
containerIDs []ids.ID,
nodeID ids.NodeID,
) InboundMessage {
containerIDBytes := make([][]byte, len(containerIDs))
encodeIDs(containerIDs, containerIDBytes)
return &inboundMessage{
nodeID: nodeID,
op: GetAcceptedOp,
message: &p2p.GetAccepted{
ChainId: chainID[:],
RequestId: requestID,
Deadline: uint64(deadline),
ContainerIds: containerIDBytes,
},
expiration: time.Now().Add(deadline),
}
}
func InboundAccepted(
chainID ids.ID,
requestID uint32,
containerIDs []ids.ID,
nodeID ids.NodeID,
) InboundMessage {
containerIDBytes := make([][]byte, len(containerIDs))
encodeIDs(containerIDs, containerIDBytes)
return &inboundMessage{
nodeID: nodeID,
op: AcceptedOp,
message: &p2p.Accepted{
ChainId: chainID[:],
RequestId: requestID,
ContainerIds: containerIDBytes,
},
expiration: mockable.MaxTime,
}
}
func InboundPushQuery(
chainID ids.ID,
requestID uint32,
deadline time.Duration,
container []byte,
requestedHeight uint64,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: PushQueryOp,
message: &p2p.PushQuery{
ChainId: chainID[:],
RequestId: requestID,
Deadline: uint64(deadline),
Container: container,
RequestedHeight: requestedHeight,
},
expiration: time.Now().Add(deadline),
}
}
func InboundPullQuery(
chainID ids.ID,
requestID uint32,
deadline time.Duration,
containerID ids.ID,
requestedHeight uint64,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: PullQueryOp,
message: &p2p.PullQuery{
ChainId: chainID[:],
RequestId: requestID,
Deadline: uint64(deadline),
ContainerId: containerID[:],
RequestedHeight: requestedHeight,
},
expiration: time.Now().Add(deadline),
}
}
func InboundChits(
chainID ids.ID,
requestID uint32,
preferredID ids.ID,
preferredIDAtHeight ids.ID,
acceptedID ids.ID,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: QbitOp,
message: &p2p.Chits{
ChainId: chainID[:],
RequestId: requestID,
PreferredId: preferredID[:],
PreferredIdAtHeight: preferredIDAtHeight[:],
AcceptedId: acceptedID[:],
},
expiration: mockable.MaxTime,
}
}
func InboundRequest(
chainID ids.ID,
requestID uint32,
deadline time.Duration,
msg []byte,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: RequestOp,
message: &p2p.Request{
ChainId: chainID[:],
RequestId: requestID,
Deadline: uint64(deadline),
AppBytes: msg,
},
expiration: time.Now().Add(deadline),
}
}
func InboundError(
nodeID ids.NodeID,
chainID ids.ID,
requestID uint32,
errorCode int32,
errorMessage string,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: ErrorOp,
message: &p2p.Error{
ChainId: chainID[:],
RequestId: requestID,
ErrorCode: errorCode,
ErrorMessage: errorMessage,
},
expiration: mockable.MaxTime,
}
}
func InboundResponse(
chainID ids.ID,
requestID uint32,
msg []byte,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: ResponseOp,
message: &p2p.Response{
ChainId: chainID[:],
RequestId: requestID,
AppBytes: msg,
},
expiration: mockable.MaxTime,
}
}
func InboundGossip(
chainID ids.ID,
msg []byte,
nodeID ids.NodeID,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: GossipOp,
message: &p2p.Gossip{
ChainId: chainID[:],
AppBytes: msg,
},
expiration: mockable.MaxTime,
}
}
// NewInboundBFTMessage creates a new InboundMessage for bft messages.
func InboundBFTMessage(
nodeID ids.NodeID,
msg *p2p.BFT,
) InboundMessage {
return &inboundMessage{
nodeID: nodeID,
op: BFTOp,
message: msg,
expiration: mockable.MaxTime,
}
}
func encodeIDs(ids []ids.ID, result [][]byte) {
for i, id := range ids {
result[i] = id[:]
}
}