forked from didapinchegit/go_rocket_mq
-
Notifications
You must be signed in to change notification settings - Fork 15
/
Copy pathremote_cmd.go
117 lines (97 loc) · 2.39 KB
/
remote_cmd.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
package rocketmq
import (
"bytes"
"encoding/binary"
"encoding/json"
"log"
"strconv"
"sync"
)
const (
RpcType = 0
RpcOneway = 1
)
var opaque int32
var decodeLock sync.Mutex
var (
remotingVersionKey string = "rocketmq.remoting.version"
ConfigVersion int = -1
requestId int32 = 0
)
type RemotingCommand struct {
// header
Code int `json:"code"`
Language string `json:"language"`
Version int `json:"version"`
Opaque int32 `json:"opaque"`
Flag int `json:"flag"`
remark string `json:"remark"`
ExtFields interface{} `json:"extFields"`
// body
Body []byte `json:"body,omitempty"`
}
func (r *RemotingCommand) encodeHeader() []byte {
length := 4
headerData := r.buildHeader()
length += len(headerData)
if r.Body != nil {
length += len(r.Body)
}
buf := bytes.NewBuffer([]byte{})
binary.Write(buf, binary.BigEndian, length)
binary.Write(buf, binary.BigEndian, len(r.Body))
buf.Write(headerData)
return buf.Bytes()
}
func (r *RemotingCommand) buildHeader() []byte {
buf, err := json.Marshal(r)
if err != nil {
return nil
}
return buf
}
func (r *RemotingCommand) encode() []byte {
length := 4
headerData := r.buildHeader()
length += len(headerData)
if r.Body != nil {
length += len(r.Body)
}
buf := bytes.NewBuffer([]byte{})
binary.Write(buf, binary.LittleEndian, length)
binary.Write(buf, binary.LittleEndian, len(r.Body))
buf.Write(headerData)
if r.Body != nil {
buf.Write(r.Body)
}
return buf.Bytes()
}
func decodeRemoteCommand(header, body []byte) *RemotingCommand {
decodeLock.Lock()
defer decodeLock.Unlock()
cmd := &RemotingCommand{}
cmd.ExtFields = make(map[string]string)
err := json.Unmarshal(header, cmd)
if err != nil {
log.Print(err)
return nil
}
cmd.Body = body
return cmd
}
func (r *RemotingCommand) decodeCommandCustomHeader() (responseHeader SendMessageResponseHeader) {
msgId := r.ExtFields.(map[string]interface{})["msgId"].(string)
queueId, _ := strconv.Atoi(r.ExtFields.(map[string]interface{})["queueId"].(string))
queueOffset, _ := strconv.Atoi(r.ExtFields.(map[string]interface{})["queueOffset"].(string))
responseHeader = SendMessageResponseHeader{
msgId: msgId,
queueId: int32(queueId),
queueOffset: int64(queueOffset),
transactionId: "",
}
return
}
func (r *RemotingCommand) markOnewayRPC() {
bits := 1 << RpcOneway
r.Flag |= bits
}