forked from LockGit/gochat
/
server.go
156 lines (145 loc) · 4.03 KB
/
server.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
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
/**
* Created by lock
* Date: 2019-08-10
* Time: 18:32
*/
package connect
import (
"encoding/json"
"fmt"
"time"
"github.com/admpub/gochat/proto"
"github.com/admpub/gochat/tools"
"github.com/gorilla/websocket"
"github.com/sirupsen/logrus"
)
type Server struct {
Buckets []*Bucket
Options ServerOptions
bucketIdx uint32
operator Operator
}
type ServerOptions struct {
WriteWait time.Duration
PongWait time.Duration
PingPeriod time.Duration
MaxMessageSize int64
ReadBufferSize int
WriteBufferSize int
BroadcastSize int
}
func NewServer(b []*Bucket, o Operator, options ServerOptions) *Server {
s := new(Server)
s.Buckets = b
s.Options = options
s.bucketIdx = uint32(len(b))
s.operator = o
return s
}
//reduce lock competition, use google city hash insert to different bucket
func (s *Server) Bucket(userId int) *Bucket {
userIdStr := fmt.Sprintf("%d", userId)
idx := tools.CityHash32([]byte(userIdStr), uint32(len(userIdStr))) % s.bucketIdx
return s.Buckets[idx]
}
func (s *Server) writePump(ch *Channel, c *Connect) {
//PingPeriod default eq 54s
ticker := time.NewTicker(s.Options.PingPeriod)
defer func() {
ticker.Stop()
ch.conn.Close()
}()
for {
select {
case message, ok := <-ch.broadcast:
//write data dead time , like http timeout , default 10s
ch.conn.SetWriteDeadline(time.Now().Add(s.Options.WriteWait))
if !ok {
logrus.Warn("SetWriteDeadline not ok")
ch.conn.WriteMessage(websocket.CloseMessage, []byte{})
return
}
w, err := ch.conn.NextWriter(websocket.TextMessage)
if err != nil {
logrus.Warn(" ch.conn.NextWriter err :%s ", err.Error())
return
}
logrus.Infof("message write body:%s", message.Body)
w.Write(message.Body)
if err := w.Close(); err != nil {
return
}
case <-ticker.C:
//heartbeat,if ping error will exit and close current websocket conn
ch.conn.SetWriteDeadline(time.Now().Add(s.Options.WriteWait))
logrus.Infof("websocket.PingMessage :%v", websocket.PingMessage)
if err := ch.conn.WriteMessage(websocket.PingMessage, nil); err != nil {
return
}
}
}
}
func (s *Server) readPump(ch *Channel, c *Connect) {
defer func() {
logrus.Infof("start exec disConnect ...")
if ch.Room == nil || ch.userId == 0 {
logrus.Infof("roomId and userId eq 0")
ch.conn.Close()
return
}
logrus.Infof("exec disConnect ...")
disConnectRequest := new(proto.DisConnectRequest)
disConnectRequest.RoomId = ch.Room.Id
disConnectRequest.UserId = ch.userId
s.Bucket(ch.userId).DeleteChannel(ch)
if err := s.operator.DisConnect(disConnectRequest); err != nil {
logrus.Warnf("DisConnect err :%s", err.Error())
}
ch.conn.Close()
}()
ch.conn.SetReadLimit(s.Options.MaxMessageSize)
ch.conn.SetReadDeadline(time.Now().Add(s.Options.PongWait))
ch.conn.SetPongHandler(func(string) error {
ch.conn.SetReadDeadline(time.Now().Add(s.Options.PongWait))
return nil
})
for {
_, message, err := ch.conn.ReadMessage()
if err != nil {
if websocket.IsUnexpectedCloseError(err, websocket.CloseGoingAway, websocket.CloseAbnormalClosure) {
logrus.Errorf("readPump ReadMessage err:%s", err.Error())
return
}
}
if message == nil {
return
}
var connReq *proto.ConnectRequest
logrus.Infof("get a message :%s", message)
if err := json.Unmarshal([]byte(message), &connReq); err != nil {
logrus.Errorf("message struct %+v", connReq)
}
if connReq.AuthToken == "" {
logrus.Errorf("s.operator.Connect no authToken")
return
}
connReq.ServerId = c.ServerId //config.Conf.Connect.ConnectWebsocket.ServerId
userId, err := s.operator.Connect(connReq)
if err != nil {
logrus.Errorf("s.operator.Connect error %s", err.Error())
return
}
if userId == 0 {
logrus.Error("Invalid AuthToken ,userId empty")
return
}
logrus.Infof("websocket rpc call return userId:%d,RoomId:%d", userId, connReq.RoomId)
b := s.Bucket(userId)
//insert into a bucket
err = b.Put(userId, connReq.RoomId, ch)
if err != nil {
logrus.Errorf("conn close err: %s", err.Error())
ch.conn.Close()
}
}
}