/
websocket_transport.go
executable file
·106 lines (97 loc) · 2.3 KB
/
websocket_transport.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
package transport
import (
"errors"
"net"
"sync"
"github.com/gearghost/emission"
"github.com/gearghost/go-protoo/logger"
"github.com/gorilla/websocket"
)
type WebSocketTransport struct {
emission.Emitter
socket *websocket.Conn
mutex *sync.Mutex
closed bool
}
func NewWebSocketTransport(socket *websocket.Conn) *WebSocketTransport {
var transport WebSocketTransport
transport.Emitter = *emission.NewEmitter()
transport.socket = socket
transport.mutex = new(sync.Mutex)
transport.closed = false
transport.socket.SetCloseHandler(func(code int, text string) error {
//logger.Warnf("%s [%d]", text, code)
transport.Emit("close", code, text)
transport.closed = true
return nil
})
return &transport
}
func (transport *WebSocketTransport) ReadMessage() {
in := make(chan []byte)
stop := make(chan struct{})
//pingTicker := time.NewTicker(pingPeriod)
var c = transport.socket
go func() {
for {
_, message, err := c.ReadMessage()
if err != nil {
logger.Warnf("Got error: %v", err)
if c, k := err.(*websocket.CloseError); k {
transport.Emit("error", c.Code, c.Text)
} else {
if c, k := err.(*net.OpError); k {
transport.Emit("error", 1008, c.Error())
}
}
close(stop)
break
}
in <- message
}
}()
for {
select {
/*case _ = <-pingTicker.C:
logger.Infof("Send keepalive !!!")
if err := transport.Send("{}"); err != nil {
logger.Errorf("Keepalive has failed")
pingTicker.Stop()
return
}*/
case message := <-in:
{
logger.Infof("Recivied data: %s", message)
transport.Emit("message", []byte(message))
}
case <-stop:
return
}
}
}
/*
* Send |message| to the connection.
*/
func (transport *WebSocketTransport) Send(message string) error {
logger.Infof("Send data: %s", message)
transport.mutex.Lock()
defer transport.mutex.Unlock()
if transport.closed {
return errors.New("websocket: write closed")
}
return transport.socket.WriteMessage(websocket.TextMessage, []byte(message))
}
/*
* Close connection.
*/
func (transport *WebSocketTransport) Close() {
transport.mutex.Lock()
defer transport.mutex.Unlock()
if transport.closed == false {
logger.Infof("Close ws transport now : ", transport)
transport.socket.Close()
transport.closed = true
} else {
logger.Warnf("Transport already closed :", transport)
}
}