mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-08-28 02:44:59 +00:00
parent
dd4685d43b
commit
8cd53bf2c2
3 changed files with 7 additions and 27 deletions
12
broadcast.go
12
broadcast.go
|
|
@ -12,8 +12,6 @@
|
|||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
//go:generate stringer -type=msgType
|
||||
|
||||
package pilosa
|
||||
|
||||
import (
|
||||
|
|
@ -55,7 +53,7 @@ func (nopBroadcaster) SendTo(*Node, Message) error { return nil }
|
|||
|
||||
// Broadcast message types.
|
||||
const (
|
||||
messageTypeCreateShard msgType = iota
|
||||
messageTypeCreateShard = iota
|
||||
messageTypeCreateIndex
|
||||
messageTypeDeleteIndex
|
||||
messageTypeCreateField
|
||||
|
|
@ -73,8 +71,6 @@ const (
|
|||
messageTypeNodeStatus
|
||||
)
|
||||
|
||||
type msgType byte
|
||||
|
||||
// MarshalInternalMessage serializes the pilosa message and adds pilosa internal
|
||||
// type info which is used by the internal messaging stuff.
|
||||
func MarshalInternalMessage(m Message, s Serializer) ([]byte, error) {
|
||||
|
|
@ -83,11 +79,11 @@ func MarshalInternalMessage(m Message, s Serializer) ([]byte, error) {
|
|||
if err != nil {
|
||||
return nil, errors.Wrap(err, "marshaling")
|
||||
}
|
||||
return append([]byte{byte(typ)}, buf...), nil
|
||||
return append([]byte{typ}, buf...), nil
|
||||
}
|
||||
|
||||
func getMessage(typ byte) Message {
|
||||
switch msgType(typ) {
|
||||
switch typ {
|
||||
case messageTypeCreateShard:
|
||||
return &CreateShardMessage{}
|
||||
case messageTypeCreateIndex:
|
||||
|
|
@ -125,7 +121,7 @@ func getMessage(typ byte) Message {
|
|||
}
|
||||
}
|
||||
|
||||
func getMessageType(m Message) msgType {
|
||||
func getMessageType(m Message) byte {
|
||||
switch m.(type) {
|
||||
case *CreateShardMessage:
|
||||
return messageTypeCreateShard
|
||||
|
|
|
|||
|
|
@ -1,16 +0,0 @@
|
|||
// Code generated by "stringer -type=msgType"; DO NOT EDIT.
|
||||
|
||||
package pilosa
|
||||
|
||||
import "strconv"
|
||||
|
||||
const _msgType_name = "messageTypeCreateShardmessageTypeCreateIndexmessageTypeDeleteIndexmessageTypeCreateFieldmessageTypeDeleteFieldmessageTypeCreateViewmessageTypeDeleteViewmessageTypeClusterStatusmessageTypeResizeInstructionmessageTypeResizeInstructionCompletemessageTypeSetCoordinatormessageTypeUpdateCoordinatormessageTypeNodeStatemessageTypeRecalculateCachesmessageTypeNodeEventmessageTypeNodeStatus"
|
||||
|
||||
var _msgType_index = [...]uint16{0, 22, 44, 66, 88, 110, 131, 152, 176, 204, 240, 265, 293, 313, 341, 361, 382}
|
||||
|
||||
func (i msgType) String() string {
|
||||
if i >= msgType(len(_msgType_index)-1) {
|
||||
return "msgType(" + strconv.FormatInt(int64(i), 10) + ")"
|
||||
}
|
||||
return _msgType_name[_msgType_index[i]:_msgType_index[i+1]]
|
||||
}
|
||||
|
|
@ -580,7 +580,7 @@ func (s *Server) SendSync(m Message) error {
|
|||
if err != nil {
|
||||
return fmt.Errorf("marshaling message: %v", err)
|
||||
}
|
||||
msg = append([]byte{byte(getMessageType(m))}, msg...)
|
||||
msg = append([]byte{getMessageType(m)}, msg...)
|
||||
for _, node := range s.cluster.nodes {
|
||||
node := node
|
||||
s.logger.Printf("SendSync to: %s", node.URI)
|
||||
|
|
@ -604,12 +604,12 @@ func (s *Server) SendAsync(m Message) error {
|
|||
|
||||
// SendTo represents an implementation of Broadcaster.
|
||||
func (s *Server) SendTo(to *Node, m Message) error {
|
||||
s.logger.Printf("SendTo: %s, type: %s", to.URI, getMessageType(m))
|
||||
s.logger.Printf("SendTo: %s", to.URI)
|
||||
msg, err := s.serializer.Marshal(m)
|
||||
if err != nil {
|
||||
return fmt.Errorf("marshaling message: %v", err)
|
||||
}
|
||||
msg = append([]byte{byte(getMessageType(m))}, msg...)
|
||||
msg = append([]byte{getMessageType(m)}, msg...)
|
||||
return s.defaultClient.SendMessage(context.Background(), &to.URI, msg)
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue