коммит произвёл
GitHub
родитель
cf3ba6661d
Коммит
462a21dcd1
@@ -34,6 +34,8 @@ const avgReadMsgSizeBytes = 1024
|
||||
|
||||
// WebSocketClient stores the necessary information required to
|
||||
// communicate with a WebSocket endpoint.
|
||||
// A client must read from PingTimeoutChannel, EventChannel and ResponseChannel to prevent
|
||||
// deadlocks from occuring in the program.
|
||||
type WebSocketClient struct {
|
||||
Url string // The location of the server like "ws://localhost:8065"
|
||||
ApiUrl string // The API location of the server like "ws://localhost:8065/api/v3"
|
||||
@@ -51,6 +53,7 @@ type WebSocketClient struct {
|
||||
quitPingWatchdog chan struct{}
|
||||
|
||||
quitWriterChan chan struct{}
|
||||
resetTimerChan chan struct{}
|
||||
closed int32
|
||||
}
|
||||
|
||||
@@ -81,6 +84,7 @@ func NewWebSocketClientWithDialer(dialer *websocket.Dialer, url, authToken strin
|
||||
writeChan: make(chan writeMessage),
|
||||
quitPingWatchdog: make(chan struct{}),
|
||||
quitWriterChan: make(chan struct{}),
|
||||
resetTimerChan: make(chan struct{}),
|
||||
}
|
||||
|
||||
client.configurePingHandling()
|
||||
@@ -125,6 +129,7 @@ func (wsc *WebSocketClient) ConnectWithDialer(dialer *websocket.Dialer) *AppErro
|
||||
wsc.writeChan = make(chan writeMessage)
|
||||
wsc.quitWriterChan = make(chan struct{})
|
||||
go wsc.writer()
|
||||
wsc.resetTimerChan = make(chan struct{})
|
||||
wsc.quitPingWatchdog = make(chan struct{})
|
||||
}
|
||||
|
||||
@@ -186,6 +191,7 @@ func (wsc *WebSocketClient) Listen() {
|
||||
close(wsc.EventChannel)
|
||||
close(wsc.ResponseChannel)
|
||||
close(wsc.quitPingWatchdog)
|
||||
close(wsc.resetTimerChan)
|
||||
// We CAS to 1 and proceed.
|
||||
if !atomic.CompareAndSwapInt32(&wsc.closed, 0, 1) {
|
||||
return
|
||||
@@ -284,21 +290,32 @@ func (wsc *WebSocketClient) pingHandler(appData string) error {
|
||||
if atomic.LoadInt32(&wsc.closed) == 1 {
|
||||
return nil
|
||||
}
|
||||
if !wsc.pingTimeoutTimer.Stop() {
|
||||
<-wsc.pingTimeoutTimer.C
|
||||
}
|
||||
|
||||
wsc.pingTimeoutTimer.Reset(time.Second * (60 + PING_TIMEOUT_BUFFER_SECONDS))
|
||||
wsc.resetTimerChan <- struct{}{}
|
||||
wsc.writeChan <- writeMessage{
|
||||
msgType: msgTypePong,
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// pingWatchdog is used to send values to the PingTimeoutChannel whenever a timeout occurs.
|
||||
// We use the resetTimerChan from the pingHandler to pass the signal, and then reset the timer
|
||||
// after draining it. And if the timer naturally expires, we also extend it to prevent it from
|
||||
// being deadlocked when the resetTimerChan case runs. Because timer.Stop would return false,
|
||||
// and the code would be forever stuck trying to read from C.
|
||||
func (wsc *WebSocketClient) pingWatchdog() {
|
||||
select {
|
||||
case <-wsc.pingTimeoutTimer.C:
|
||||
wsc.PingTimeoutChannel <- true
|
||||
case <-wsc.quitPingWatchdog:
|
||||
for {
|
||||
select {
|
||||
case <-wsc.resetTimerChan:
|
||||
if !wsc.pingTimeoutTimer.Stop() {
|
||||
<-wsc.pingTimeoutTimer.C
|
||||
}
|
||||
wsc.pingTimeoutTimer.Reset(time.Second * (60 + PING_TIMEOUT_BUFFER_SECONDS))
|
||||
|
||||
case <-wsc.pingTimeoutTimer.C:
|
||||
wsc.PingTimeoutChannel <- true
|
||||
wsc.pingTimeoutTimer.Reset(time.Second * (60 + PING_TIMEOUT_BUFFER_SECONDS))
|
||||
case <-wsc.quitPingWatchdog:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Ссылка в новой задаче
Block a user