Add Ping Timeout Signal for WebSocket Client (#8666)
Import Order gofmt gofmt gofmt Renamed Const to make Unit more obvious Renamed Const to make Unit more obvious #2
Этот коммит содержится в:
коммит произвёл
Joram Wilander
родитель
d8dd271e43
Коммит
7fa1c6c4ba
@@ -6,12 +6,14 @@ package model
|
|||||||
import (
|
import (
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"time"
|
||||||
|
|
||||||
"github.com/gorilla/websocket"
|
"github.com/gorilla/websocket"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
SOCKET_MAX_MESSAGE_SIZE_KB = 8 * 1024 // 8KB
|
SOCKET_MAX_MESSAGE_SIZE_KB = 8 * 1024 // 8KB
|
||||||
|
PING_TIMEOUT_BUFFER_SECONDS = 5
|
||||||
)
|
)
|
||||||
|
|
||||||
type WebSocketClient struct {
|
type WebSocketClient struct {
|
||||||
@@ -21,9 +23,11 @@ type WebSocketClient struct {
|
|||||||
Conn *websocket.Conn // The WebSocket connection
|
Conn *websocket.Conn // The WebSocket connection
|
||||||
AuthToken string // The token used to open the WebSocket
|
AuthToken string // The token used to open the WebSocket
|
||||||
Sequence int64 // The ever-incrementing sequence attached to each WebSocket action
|
Sequence int64 // The ever-incrementing sequence attached to each WebSocket action
|
||||||
|
PingTimeoutChannel chan bool // The channel used to signal ping timeouts
|
||||||
EventChannel chan *WebSocketEvent
|
EventChannel chan *WebSocketEvent
|
||||||
ResponseChannel chan *WebSocketResponse
|
ResponseChannel chan *WebSocketResponse
|
||||||
ListenError *AppError
|
ListenError *AppError
|
||||||
|
pingTimeoutTimer *time.Timer
|
||||||
}
|
}
|
||||||
|
|
||||||
// NewWebSocketClient constructs a new WebSocket client with convenience
|
// NewWebSocketClient constructs a new WebSocket client with convenience
|
||||||
@@ -47,11 +51,15 @@ func NewWebSocketClientWithDialer(dialer *websocket.Dialer, url, authToken strin
|
|||||||
conn,
|
conn,
|
||||||
authToken,
|
authToken,
|
||||||
1,
|
1,
|
||||||
|
make(chan bool, 1),
|
||||||
make(chan *WebSocketEvent, 100),
|
make(chan *WebSocketEvent, 100),
|
||||||
make(chan *WebSocketResponse, 100),
|
make(chan *WebSocketResponse, 100),
|
||||||
nil,
|
nil,
|
||||||
|
nil,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
client.configurePingHandling()
|
||||||
|
|
||||||
client.SendMessage(WEBSOCKET_AUTHENTICATION_CHALLENGE, map[string]interface{}{"token": authToken})
|
client.SendMessage(WEBSOCKET_AUTHENTICATION_CHALLENGE, map[string]interface{}{"token": authToken})
|
||||||
|
|
||||||
return client, nil
|
return client, nil
|
||||||
@@ -78,11 +86,15 @@ func NewWebSocketClient4WithDialer(dialer *websocket.Dialer, url, authToken stri
|
|||||||
conn,
|
conn,
|
||||||
authToken,
|
authToken,
|
||||||
1,
|
1,
|
||||||
|
make(chan bool, 1),
|
||||||
make(chan *WebSocketEvent, 100),
|
make(chan *WebSocketEvent, 100),
|
||||||
make(chan *WebSocketResponse, 100),
|
make(chan *WebSocketResponse, 100),
|
||||||
nil,
|
nil,
|
||||||
|
nil,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
client.configurePingHandling()
|
||||||
|
|
||||||
client.SendMessage(WEBSOCKET_AUTHENTICATION_CHALLENGE, map[string]interface{}{"token": authToken})
|
client.SendMessage(WEBSOCKET_AUTHENTICATION_CHALLENGE, map[string]interface{}{"token": authToken})
|
||||||
|
|
||||||
return client, nil
|
return client, nil
|
||||||
@@ -99,6 +111,8 @@ func (wsc *WebSocketClient) ConnectWithDialer(dialer *websocket.Dialer) *AppErro
|
|||||||
return NewAppError("Connect", "model.websocket_client.connect_fail.app_error", nil, err.Error(), http.StatusInternalServerError)
|
return NewAppError("Connect", "model.websocket_client.connect_fail.app_error", nil, err.Error(), http.StatusInternalServerError)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
wsc.configurePingHandling()
|
||||||
|
|
||||||
wsc.EventChannel = make(chan *WebSocketEvent, 100)
|
wsc.EventChannel = make(chan *WebSocketEvent, 100)
|
||||||
wsc.ResponseChannel = make(chan *WebSocketResponse, 100)
|
wsc.ResponseChannel = make(chan *WebSocketResponse, 100)
|
||||||
|
|
||||||
@@ -181,3 +195,24 @@ func (wsc *WebSocketClient) GetStatusesByIds(userIds []string) {
|
|||||||
}
|
}
|
||||||
wsc.SendMessage("get_statuses_by_ids", data)
|
wsc.SendMessage("get_statuses_by_ids", data)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (wsc *WebSocketClient) configurePingHandling() {
|
||||||
|
wsc.Conn.SetPingHandler(wsc.pingHandler)
|
||||||
|
wsc.pingTimeoutTimer = time.NewTimer(time.Second * (60 + PING_TIMEOUT_BUFFER_SECONDS))
|
||||||
|
go wsc.pingWatchdog()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (wsc *WebSocketClient) pingHandler(appData string) error {
|
||||||
|
if !wsc.pingTimeoutTimer.Stop() {
|
||||||
|
<-wsc.pingTimeoutTimer.C
|
||||||
|
}
|
||||||
|
|
||||||
|
wsc.pingTimeoutTimer.Reset(time.Second * (60 + PING_TIMEOUT_BUFFER_SECONDS))
|
||||||
|
wsc.Conn.WriteMessage(websocket.PongMessage, []byte{})
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (wsc *WebSocketClient) pingWatchdog() {
|
||||||
|
<-wsc.pingTimeoutTimer.C
|
||||||
|
wsc.PingTimeoutChannel <- true
|
||||||
|
}
|
||||||
|
|||||||
Ссылка в новой задаче
Block a user