-
Notifications
You must be signed in to change notification settings - Fork 2
/
http.go
214 lines (195 loc) · 4.85 KB
/
http.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
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
package xhttp
import (
"encoding/json"
"errors"
"fmt"
"io/ioutil"
"net/http"
"sync"
"sync/atomic"
"time"
"github.com/eyasliu/cs"
)
// HTTP cs 的 HTTP 适配器
type HTTP struct {
srv *cs.Srv
receive chan *reqMessage
session map[string][]*SSEConn // http 模式可能出现一个会话多个连接的情况
sessionMu sync.RWMutex
sidKey string
sidCount uint32
hbTime time.Duration
msgType SSEMsgType
}
var defaultHeartBeatTime = 10 * time.Second
// New 实例化适配器
func New() *HTTP {
h := &HTTP{
sidKey: "sid",
session: make(map[string][]*SSEConn),
receive: make(chan *reqMessage, 2),
hbTime: defaultHeartBeatTime,
msgType: SSEMessage,
}
return h
}
// Handler impl http.HandlerFunc to handler http request
func (h *HTTP) Handler(w http.ResponseWriter, req *http.Request) {
if h.srv == nil {
w.WriteHeader(500)
w.Write([]byte("srv not running"))
return
}
sid := h.setSid(w, req)
if sid == "" {
w.WriteHeader(400)
w.Write([]byte("invalid sid, must allow cookie to store sid"))
return
}
if req.Method == "GET" {
h.invokeSSE(sid, w, req)
} else if req.Method == "POST" || req.Method == "PUT" || req.Method == "DELETE" {
h.invokeHandle(sid, w, req)
}
}
// ServeHTTP impl http.Handler to handler http request
func (h *HTTP) ServeHTTP(w http.ResponseWriter, req *http.Request) {
h.Handler(w, req)
}
// Srv 返回cs.Srv 实例,如果没有绑定实例则初始化一个
func (h *HTTP) Srv() *cs.Srv {
if h.srv == nil {
h.srv = cs.New(h)
}
return h.srv
}
// Write 实现 cs.ServerAdapter 接口,给连接推送消息
func (h *HTTP) Write(sid string, resp *cs.Response) error {
conns, ok := h.session[sid]
if !ok {
return errors.New("connection is already close")
}
for _, conn := range conns {
if err := conn.Send(resp); err != nil {
return err
}
}
return nil
}
// Read 实现 cs.ServerAdapter 接口,读取消息,每次返回一条,循环读取
func (h *HTTP) Read(srv *cs.Srv) (sid string, req *cs.Request, err error) {
h.srv = srv
<-make(chan struct{})
return "", nil, errors.New("HTTP Adapter unsupport Read")
}
// Close 实现 cs.ServerAdapter 接口,关闭指定连接
func (h *HTTP) Close(sid string) error {
conns, ok := h.session[sid]
if !ok {
return errors.New("ths sid already close")
}
for _, conn := range conns {
conn.destroy(nil)
}
h.sessionMu.Lock()
delete(h.session, sid)
h.sessionMu.Unlock()
return nil
}
// GetAllSID 实现 cs.ServerAdapter 接口,获取当前服务所有SID,用于遍历连接
func (h *HTTP) GetAllSID() []string {
sids := make([]string, 0, len(h.session))
h.sessionMu.RLock()
for sid := range h.session {
sids = append(sids, sid)
}
h.sessionMu.RUnlock()
return sids
}
// 基于cookie,设置会话sid
func (h *HTTP) setSid(w http.ResponseWriter, req *http.Request) string {
cookie, err := req.Cookie(h.sidKey)
var sid string
if err != nil || cookie == nil {
atomic.AddUint32(&h.sidCount, 1)
// 因为sid是存cookie的,而程序每次重启,这个计数器都会重置为 0
// 只使用计数器会导致 sid 重复,需要加上其他变量,计数器可以保证在高并发时不会重复
sid = fmt.Sprintf("http.%d-%d", time.Now().Unix(), h.sidCount)
cookie = &http.Cookie{
Name: h.sidKey,
Value: sid,
HttpOnly: true,
// Expires: time.Now().Add(24 * time.Hour)
}
http.SetCookie(w, cookie)
} else {
sid = cookie.Value
}
return sid
}
// 处理cmd路由
func (h *HTTP) invokeHandle(sid string, w http.ResponseWriter, req *http.Request) {
reqData := &requestData{}
respData := &cs.Response{}
data, err := ioutil.ReadAll(req.Body)
if err != nil {
respData.Msg = err.Error()
} else {
err := json.Unmarshal(data, reqData)
if err != nil {
respData.Msg = err.Error()
}
ctx := h.srv.NewContext(h, sid, &cs.Request{
Cmd: reqData.Cmd,
Seqno: reqData.Seqno,
RawData: reqData.Data,
})
h.srv.CallContext(ctx)
respData = ctx.Response
}
resp := &responseData{
Cmd: respData.Cmd,
Seqno: respData.Seqno,
Code: respData.Code,
Msg: respData.Msg,
Data: respData.Data,
}
respBt, _ := json.Marshal(resp)
w.WriteHeader(200)
w.Header().Set("Content-Type", "application/json")
w.Write(respBt)
}
// 处理 sse 连接
func (h *HTTP) invokeSSE(sid string, w http.ResponseWriter, req *http.Request) {
conn, err := newSSEConn(w, h.msgType, h.hbTime)
if err != nil {
return
}
h.sessionMu.Lock()
conns, ok := h.session[sid]
if !ok {
conns = []*SSEConn{conn}
} else {
conns = append(conns, conn)
}
h.session[sid] = conns
h.sessionMu.Unlock()
<-conn.notifyErr
h.sessionMu.Lock()
conns, ok = h.session[sid]
if !ok {
return
}
nextConns := make([]*SSEConn, 0, len(conns)-1)
for _, c := range conns {
if c != conn {
nextConns = append(nextConns, c)
}
}
if len(nextConns) > 0 {
h.session[sid] = nextConns
} else {
delete(h.session, sid)
}
h.sessionMu.Unlock()
}