-
Notifications
You must be signed in to change notification settings - Fork 3
/
typeProcessor.go
55 lines (47 loc) · 1.21 KB
/
typeProcessor.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
package sseKit
import (
"github.com/gin-gonic/gin"
"github.com/richelieu-yang/chimera/v2/src/component/web/push/pushKit"
"github.com/richelieu-yang/chimera/v2/src/mutexKit"
"net/http"
)
type SseProcessor struct {
msgType messageType
listeners pushKit.Listeners
}
func (p *SseProcessor) HandleWithGin(ctx *gin.Context) {
p.Handle(ctx.Writer, ctx.Request)
}
func (p *SseProcessor) Handle(w http.ResponseWriter, r *http.Request) {
if errText := IsSseSupported(w, r); errText != "" {
p.listeners.OnFailure(w, r, errText)
return
}
// 设置 response header
SetHeaders(w)
channel := p.newChannel(w, r)
p.listeners.OnHandshake(w, r, channel)
select {
case <-r.Context().Done():
p.listeners.OnClose(channel, "Context done")
case <-w.(http.CloseNotifier).CloseNotify():
// SSE客户端关闭后,会走此处
p.listeners.OnClose(channel, "Connection closed")
}
}
func (p *SseProcessor) newChannel(w http.ResponseWriter, r *http.Request) pushKit.Channel {
return &SseChannel{
BaseChannel: &pushKit.BaseChannel{
Id: "",
Bsid: "",
User: "",
Group: "",
RWMutex: mutexKit.RWMutex{},
Data: nil,
Closed: false,
Listeners: p.listeners,
},
w: w,
r: r,
}
}