forked from DrmagicE/gmqtt
/
ws.go
82 lines (72 loc) · 1.68 KB
/
ws.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
package mqtt
import (
"net"
"net/http"
"github.com/gorilla/websocket"
"go.uber.org/zap"
)
// WsServer is used to build websocket server
type WsServer struct {
Server *http.Server
Path string // Url path
CertFile string //TLS configration
KeyFile string //TLS configration
}
var defaultUpgrader = &websocket.Upgrader{
ReadBufferSize: readBufferSize,
WriteBufferSize: writeBufferSize,
CheckOrigin: func(r *http.Request) bool {
return true
},
Subprotocols: []string{"mqtt"},
}
//实现io.ReadWriter接口
// wsConn implements the io.ReadWriter
type wsConn struct {
net.Conn
c *websocket.Conn
}
func (ws *wsConn) Close() error {
return ws.Conn.Close()
}
func (ws *wsConn) Read(p []byte) (n int, err error) {
msgType, r, err := ws.c.NextReader()
if err != nil {
return 0, err
}
if msgType != websocket.BinaryMessage {
return 0, ErrInvalWsMsgType
}
return r.Read(p)
}
func (ws *wsConn) Write(p []byte) (n int, err error) {
err = ws.c.WriteMessage(websocket.BinaryMessage, p)
if err != nil {
return 0, err
}
return len(p), err
}
func (srv *server) serveWebSocket(ws *WsServer) {
var err error
if ws.CertFile != "" && ws.KeyFile != "" {
err = ws.Server.ListenAndServeTLS(ws.CertFile, ws.KeyFile)
} else {
err = ws.Server.ListenAndServe()
}
if err != http.ErrServerClosed {
panic(err.Error())
}
}
func (srv *server) wsHandler() http.HandlerFunc {
return func(w http.ResponseWriter, r *http.Request) {
c, err := defaultUpgrader.Upgrade(w, r, nil)
if err != nil {
zaplog.Warn("websocket upgrade error", zap.String("msg", err.Error()))
return
}
defer c.Close()
conn := &wsConn{c.UnderlyingConn(), c}
client := srv.newClient(conn)
client.serve()
}
}