-
Notifications
You must be signed in to change notification settings - Fork 0
/
readwrite.go
131 lines (106 loc) · 2.98 KB
/
readwrite.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
118
119
120
121
122
123
124
125
126
127
128
129
130
131
package ws
// #
// # NOT used
// # This should be deleted - July 2017
import (
// "bytes"
// "encoding/gob"
"github.com/gorilla/websocket"
log "github.com/sirupsen/logrus"
"time"
)
func channelIsClosed(ch <-chan WSMsg) bool {
select {
case <-ch:
return true
default:
}
return false
}
func (wsConn *WSConn) WaitResponse(timeout time.Duration) (msg WSMsg, err error) {
// if channelIsClosed(wsConn.readChannel) {
// log.Errorf("WS wait channel is closed")
// return WSMsg{}, ErrorInternal
// }
select {
case msg := <-wsConn.readChannel:
return msg, nil
case <-time.After(timeout):
log.WithFields(log.Fields{"clientId": wsConn.ClientId}).Error("WS wait timed out")
// close(wsConn.readChannel)
wsConn.c.Close() // force readerLoop to exit and redial
return WSMsg{}, ErrorTimeout
// default:
}
}
func (wsConn *WSConn) read() (msg WSMsg, err error) {
// log.Info("read")
wsConn.readMutex.Lock()
_, packet, err := wsConn.c.ReadMessage()
wsConn.readMutex.Unlock()
if err != nil {
if e, ok := err.(*websocket.CloseError); ok && e.Code == websocket.CloseNormalClosure {
return nil, ErrorClosed
}
return nil, err
}
if len(packet) < 2 {
log.WithFields(log.Fields{"clientId": wsConn.ClientId}).Error("WS invalid msg", packet)
return nil, ErrorInvalidRequest
}
msg = WSMsg(packet)
return msg, nil
}
func (wsConn *WSConn) Write(msg WSMsg) (seq uint8, err error) {
// log.Info("write")
// Add sequence number if this is new message
// if msg.Sequence() == msgSeqNew {
if !msg.IsResponse() {
wsConn.msgSeq = (wsConn.msgSeq + 1) & 0x7f // drop the first bit
msg.setSequence(wsConn.msgSeq)
}
wsConn.writeMutex.Lock()
defer wsConn.writeMutex.Unlock()
wsConn.c.SetWriteDeadline(time.Now().Add(writeWait))
if err := wsConn.c.WriteMessage(websocket.BinaryMessage, msg); err != nil {
log.WithFields(log.Fields{"clientId": wsConn.ClientId}).Error("WS failed to write websocket msg: ", err)
return 0, err
}
return msg.Sequence(), nil
}
func (wsConn *WSConn) RPC(msgType uint8, token *string, in interface{}, out interface{}) (err error) {
msg, err := Encode(msgType, token, in)
if err != nil {
return err
}
response, err := wsConn.WriteAndWaitResponse(msg)
if err != nil {
return err
}
if out != nil {
err = response.Decode(nil, out)
if err != nil {
return err
}
}
return nil
}
func (wsConn *WSConn) WriteAndWaitResponse(msg WSMsg) (response WSMsg, err error) {
// log.Info("Writing")
seq, err := wsConn.Write(msg)
if err != nil {
return WSMsg{}, err
}
// log.Info("Waiting")
response, err = wsConn.WaitResponse(readTimeout)
if err != nil {
// log.WithFields(log.Fields{"clientId": wsConn.ClientId}).Error("WS waiting response ", err)
return WSMsg{}, err
}
if seq != response.Sequence() {
log.WithFields(log.Fields{"clientId": wsConn.ClientId}).Error("WS unexpected sequence: ")
return WSMsg{}, ErrorInternal
}
// log.WithFields(log.Fields{"clientId": wsConn.ClientId}).Info("WS received response ", response)
return response, nil
}