-
Notifications
You must be signed in to change notification settings - Fork 0
/
Copy pathmessage.go
119 lines (112 loc) · 2.22 KB
/
message.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
package websocket
import (
"context"
"errors"
"io"
)
type Message struct {
io.Reader
OpCode OpCode
}
func (w *webSocket) sendMessage(message *Message) error {
ctx := context.Background()
frame := &Frame{
Payload: nil,
Fin: false,
Mask: w.mask,
OpCode: message.OpCode,
}
buf := make([]byte, 2048)
offset := 0
if message.Reader == nil {
message.Reader = emptyReader
}
for {
n, err := message.Read(buf[offset:])
if err != nil && err != io.EOF {
return err
}
offset += n
if err == nil && n < len(buf) {
continue
}
frame.Payload = &io.LimitedReader{
R: newBytesBuffer(buf[:offset]),
N: int64(offset),
}
frame.Fin = err != nil
err = w.sendFrame(ctx, frame)
if err != nil {
return err
}
if frame.Fin {
return nil
}
offset = 0
frame.OpCode = ContinuationFrame
}
}
func (w *webSocket) SendMessage(message *Message) error {
w.sendLock.Lock()
defer w.sendLock.Unlock()
return w.sendMessage(message)
}
var ErrPreviousMessageNotReadToCompletion = errors.New("previous message not read to completion")
func (w *webSocket) readMessage() (*Message, error) {
w.readLock.Lock()
ctx := context.Background()
frame, err := w.readFrame(ctx)
if err != nil {
return nil, err
}
return &Message{
Reader: rwFunc(func(b []byte) (int, error) {
for {
if frame != nil {
n, readErr := frame.Payload.Read(b)
if readErr == io.EOF && frame.Fin != true {
readErr = nil
frame = nil
}
if readErr != nil {
w.readLock.Unlock()
}
return n, readErr
}
frame, err = w.readFrame(ctx)
if err != nil {
return 0, err
}
if frame.OpCode != ContinuationFrame {
return 0, ErrPreviousMessageNotReadToCompletion
}
}
}),
OpCode: frame.OpCode,
}, nil
}
func (w *webSocket) ReadMessage() (*Message, error) {
for {
message, err := w.readMessage()
if err != nil {
return nil, err
}
if message.OpCode == Ping {
err = w.responsePong(message)
if err != nil {
return nil, err
}
} else if message.OpCode == ConnectionClose {
err = w.Close()
if err != nil {
return nil, err
}
_, err = io.Copy(blackHole, message)
if err != nil {
return nil, err
}
} else {
return message, nil
}
}
}