2020-01-09 05:24:45 +00:00
|
|
|
package wsutil
|
|
|
|
|
|
|
|
import (
|
2020-04-06 19:03:42 +00:00
|
|
|
"bytes"
|
2020-01-09 05:24:45 +00:00
|
|
|
"compress/zlib"
|
|
|
|
"context"
|
2020-01-29 03:54:22 +00:00
|
|
|
"io"
|
2020-01-09 05:24:45 +00:00
|
|
|
"net/http"
|
2020-04-11 19:34:40 +00:00
|
|
|
"sync"
|
2020-04-06 19:03:42 +00:00
|
|
|
"time"
|
2020-01-09 05:24:45 +00:00
|
|
|
|
2020-04-06 19:03:42 +00:00
|
|
|
"github.com/gorilla/websocket"
|
2020-01-09 05:24:45 +00:00
|
|
|
"github.com/pkg/errors"
|
|
|
|
)
|
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
// CopyBufferSize is used for the initial size of the internal WS' buffer. Its
|
|
|
|
// size is 4KB.
|
|
|
|
var CopyBufferSize = 4096
|
2020-04-06 19:03:42 +00:00
|
|
|
|
2020-08-20 21:15:52 +00:00
|
|
|
// MaxCapUntilReset determines the maximum capacity before the bytes buffer is
|
2020-10-28 17:19:22 +00:00
|
|
|
// re-allocated. It is roughly 16KB, quadruple CopyBufferSize.
|
|
|
|
var MaxCapUntilReset = CopyBufferSize * 4
|
2020-08-20 21:15:52 +00:00
|
|
|
|
2020-04-06 19:03:42 +00:00
|
|
|
// CloseDeadline controls the deadline to wait for sending the Close frame.
|
|
|
|
var CloseDeadline = time.Second
|
2020-01-09 05:24:45 +00:00
|
|
|
|
2020-04-11 19:34:40 +00:00
|
|
|
// ErrWebsocketClosed is returned if the websocket is already closed.
|
2020-05-16 21:14:49 +00:00
|
|
|
var ErrWebsocketClosed = errors.New("websocket is closed")
|
2020-04-11 19:34:40 +00:00
|
|
|
|
2020-01-09 05:24:45 +00:00
|
|
|
// Connection is an interface that abstracts around a generic Websocket driver.
|
2020-04-06 19:03:42 +00:00
|
|
|
// This connection expects the driver to handle compression by itself, including
|
|
|
|
// modifying the connection URL.
|
2020-01-09 05:24:45 +00:00
|
|
|
type Connection interface {
|
2020-01-15 04:43:34 +00:00
|
|
|
// Dial dials the address (string). Context needs to be passed in for
|
|
|
|
// timeout. This method should also be re-usable after Close is called.
|
|
|
|
Dial(context.Context, string) error
|
|
|
|
|
2020-01-09 05:24:45 +00:00
|
|
|
// Listen sends over events constantly. Error will be non-nil if Data is
|
|
|
|
// nil, so check for Error first.
|
|
|
|
Listen() <-chan Event
|
|
|
|
|
2020-04-06 19:03:42 +00:00
|
|
|
// Send allows the caller to send bytes. Thread safety is a requirement.
|
2020-04-24 06:34:08 +00:00
|
|
|
Send(context.Context, []byte) error
|
2020-01-09 05:24:45 +00:00
|
|
|
|
|
|
|
// Close should close the websocket connection. The connection will not be
|
2020-04-11 03:03:52 +00:00
|
|
|
// reused.
|
|
|
|
Close() error
|
2020-01-09 05:24:45 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// Conn is the default Websocket connection. It compresses all payloads using
|
|
|
|
// zlib.
|
|
|
|
type Conn struct {
|
2020-10-28 17:19:22 +00:00
|
|
|
mutex sync.Mutex
|
|
|
|
|
2020-01-20 11:06:20 +00:00
|
|
|
Conn *websocket.Conn
|
2020-01-09 05:24:45 +00:00
|
|
|
|
2020-04-06 19:03:42 +00:00
|
|
|
dialer *websocket.Dialer
|
2020-01-09 05:24:45 +00:00
|
|
|
events chan Event
|
|
|
|
}
|
|
|
|
|
|
|
|
var _ Connection = (*Conn)(nil)
|
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
// NewConn creates a new default websocket connection with a default dialer.
|
2020-04-24 06:34:08 +00:00
|
|
|
func NewConn() *Conn {
|
2020-10-28 17:19:22 +00:00
|
|
|
return NewConnWithDialer(&websocket.Dialer{
|
|
|
|
Proxy: http.ProxyFromEnvironment,
|
|
|
|
HandshakeTimeout: WSTimeout,
|
|
|
|
ReadBufferSize: CopyBufferSize,
|
|
|
|
WriteBufferSize: CopyBufferSize,
|
|
|
|
EnableCompression: true,
|
|
|
|
})
|
2020-04-24 06:34:08 +00:00
|
|
|
}
|
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
// NewConn creates a new default websocket connection with a custom dialer.
|
|
|
|
func NewConnWithDialer(dialer *websocket.Dialer) *Conn {
|
|
|
|
return &Conn{dialer: dialer}
|
2020-01-09 05:24:45 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
func (c *Conn) Dial(ctx context.Context, addr string) error {
|
2020-04-06 19:03:42 +00:00
|
|
|
// Enable compression:
|
2020-10-28 17:19:22 +00:00
|
|
|
headers := http.Header{
|
|
|
|
"Accept-Encoding": {"zlib"},
|
|
|
|
}
|
2020-04-06 19:03:42 +00:00
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
// BUG which prevents stream compression.
|
|
|
|
// See https://github.com/golang/go/issues/31514.
|
2020-01-09 05:24:45 +00:00
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
conn, _, err := c.dialer.DialContext(ctx, addr, headers)
|
2020-02-11 17:23:42 +00:00
|
|
|
if err != nil {
|
2020-05-16 21:14:49 +00:00
|
|
|
return errors.Wrap(err, "failed to dial WS")
|
2020-02-11 17:23:42 +00:00
|
|
|
}
|
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
events := make(chan Event, WSBuffer)
|
|
|
|
go startReadLoop(conn, events)
|
2020-04-11 19:34:40 +00:00
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
c.mutex.Lock()
|
|
|
|
defer c.mutex.Unlock()
|
2020-04-06 21:03:08 +00:00
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
c.Conn = conn
|
|
|
|
c.events = events
|
2020-04-06 21:03:08 +00:00
|
|
|
|
2020-01-09 05:24:45 +00:00
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *Conn) Listen() <-chan Event {
|
|
|
|
return c.events
|
|
|
|
}
|
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
// resetDeadline is used to reset the write deadline after using the context's.
|
|
|
|
var resetDeadline = time.Time{}
|
|
|
|
|
|
|
|
func (c *Conn) Send(ctx context.Context, b []byte) error {
|
|
|
|
c.mutex.Lock()
|
|
|
|
defer c.mutex.Unlock()
|
|
|
|
|
|
|
|
d, ok := ctx.Deadline()
|
|
|
|
if ok {
|
|
|
|
c.Conn.SetWriteDeadline(d)
|
|
|
|
defer c.Conn.SetWriteDeadline(resetDeadline)
|
|
|
|
}
|
|
|
|
|
|
|
|
return c.Conn.WriteMessage(websocket.TextMessage, b)
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *Conn) Close() error {
|
|
|
|
// Use a sync.Once to guarantee that other Close() calls block until the
|
|
|
|
// main call is done. It also prevents future calls.
|
|
|
|
WSDebug("Conn: Acquiring write lock...")
|
|
|
|
|
|
|
|
// Acquire the write lock forever.
|
|
|
|
c.mutex.Lock()
|
|
|
|
defer c.mutex.Unlock()
|
|
|
|
|
|
|
|
WSDebug("Conn: Write lock acquired; closing.")
|
|
|
|
|
|
|
|
// Close the WS.
|
|
|
|
err := c.closeWS()
|
|
|
|
|
|
|
|
WSDebug("Conn: Websocket closed; error:", err)
|
|
|
|
WSDebug("Conn: Flusing events...")
|
|
|
|
|
|
|
|
// Flush all events before closing the channel. This will return as soon as
|
|
|
|
// c.events is closed, or after closed.
|
|
|
|
for range c.events {
|
|
|
|
}
|
|
|
|
|
|
|
|
WSDebug("Flushed events.")
|
|
|
|
|
|
|
|
// Mark c.Conn as empty.
|
|
|
|
c.Conn = nil
|
|
|
|
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
|
|
|
func (c *Conn) closeWS() error {
|
|
|
|
// We can't close with a write control here, since it will invalidate the
|
|
|
|
// old session, breaking resumes.
|
|
|
|
|
|
|
|
// // Quick deadline:
|
|
|
|
// deadline := time.Now().Add(CloseDeadline)
|
|
|
|
|
|
|
|
// // Make a closure message:
|
|
|
|
// msg := websocket.FormatCloseMessage(websocket.CloseGoingAway, "")
|
|
|
|
|
|
|
|
// // Send a close message before closing the connection. We're not error
|
|
|
|
// // checking this because it's not important.
|
|
|
|
// err = c.Conn.WriteControl(websocket.CloseMessage, msg, deadline)
|
|
|
|
|
|
|
|
return c.Conn.Close()
|
|
|
|
}
|
|
|
|
|
|
|
|
// loopState is a thread-unsafe disposable state container for the read loop.
|
|
|
|
// It's made to completely separate the read loop of any synchronization that
|
|
|
|
// doesn't involve the websocket connection itself.
|
|
|
|
type loopState struct {
|
|
|
|
conn *websocket.Conn
|
|
|
|
zlib io.ReadCloser
|
|
|
|
buf bytes.Buffer
|
|
|
|
}
|
|
|
|
|
|
|
|
func startReadLoop(conn *websocket.Conn, eventCh chan<- Event) {
|
2020-04-06 19:03:42 +00:00
|
|
|
// Clean up the events channel in the end.
|
2020-10-28 17:19:22 +00:00
|
|
|
defer close(eventCh)
|
|
|
|
|
|
|
|
// Allocate the read loop its own private resources.
|
|
|
|
state := loopState{conn: conn}
|
|
|
|
state.buf.Grow(CopyBufferSize)
|
2020-04-06 19:03:42 +00:00
|
|
|
|
|
|
|
for {
|
2020-10-28 17:19:22 +00:00
|
|
|
b, err := state.handle()
|
2020-04-06 19:03:42 +00:00
|
|
|
if err != nil {
|
|
|
|
// Is the error an EOF?
|
2020-04-09 23:19:52 +00:00
|
|
|
if errors.Is(err, io.EOF) {
|
2020-04-06 19:03:42 +00:00
|
|
|
// Yes it is, exit.
|
2020-02-02 22:12:54 +00:00
|
|
|
return
|
2020-01-09 05:24:45 +00:00
|
|
|
}
|
|
|
|
|
2020-04-06 19:03:42 +00:00
|
|
|
// Check if the error is a normal one:
|
|
|
|
if websocket.IsCloseError(err, websocket.CloseNormalClosure) {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
// Unusual error; log and exit:
|
2020-10-28 17:19:22 +00:00
|
|
|
eventCh <- Event{nil, errors.Wrap(err, "WS error")}
|
2020-04-06 19:03:42 +00:00
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2020-04-09 20:49:12 +00:00
|
|
|
// If the payload length is 0, skip it.
|
|
|
|
if len(b) == 0 {
|
2020-04-06 19:03:42 +00:00
|
|
|
continue
|
2020-01-09 05:24:45 +00:00
|
|
|
}
|
2020-04-06 19:03:42 +00:00
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
eventCh <- Event{b, nil}
|
2020-04-06 19:03:42 +00:00
|
|
|
}
|
2020-01-09 05:24:45 +00:00
|
|
|
}
|
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
func (state *loopState) handle() ([]byte, error) {
|
2020-04-06 19:03:42 +00:00
|
|
|
// skip message type
|
2020-10-28 17:19:22 +00:00
|
|
|
t, r, err := state.conn.NextReader()
|
2020-01-09 05:24:45 +00:00
|
|
|
if err != nil {
|
2020-01-16 03:28:21 +00:00
|
|
|
return nil, err
|
2020-01-09 05:24:45 +00:00
|
|
|
}
|
|
|
|
|
2020-04-06 19:03:42 +00:00
|
|
|
if t == websocket.BinaryMessage {
|
2020-01-09 05:24:45 +00:00
|
|
|
// Probably a zlib payload
|
2020-04-11 19:34:40 +00:00
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
if state.zlib == nil {
|
2020-04-11 19:34:40 +00:00
|
|
|
z, err := zlib.NewReader(r)
|
|
|
|
if err != nil {
|
2020-05-16 21:14:49 +00:00
|
|
|
return nil, errors.Wrap(err, "failed to create a zlib reader")
|
2020-04-11 19:34:40 +00:00
|
|
|
}
|
2020-10-28 17:19:22 +00:00
|
|
|
state.zlib = z
|
2020-04-11 19:34:40 +00:00
|
|
|
} else {
|
2020-10-28 17:19:22 +00:00
|
|
|
if err := state.zlib.(zlib.Resetter).Reset(r, nil); err != nil {
|
2020-05-16 21:14:49 +00:00
|
|
|
return nil, errors.Wrap(err, "failed to reset zlib reader")
|
2020-04-11 19:34:40 +00:00
|
|
|
}
|
2020-01-09 05:24:45 +00:00
|
|
|
}
|
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
defer state.zlib.Close()
|
|
|
|
r = state.zlib
|
2020-08-20 21:15:52 +00:00
|
|
|
}
|
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
return state.readAll(r)
|
2020-08-20 21:15:52 +00:00
|
|
|
}
|
|
|
|
|
2020-04-06 19:03:42 +00:00
|
|
|
// readAll reads bytes into an existing buffer, copy it over, then wipe the old
|
|
|
|
// buffer.
|
2020-10-28 17:19:22 +00:00
|
|
|
func (state *loopState) readAll(r io.Reader) ([]byte, error) {
|
|
|
|
defer state.buf.Reset()
|
2020-08-20 21:15:52 +00:00
|
|
|
|
2020-10-28 17:19:22 +00:00
|
|
|
if _, err := state.buf.ReadFrom(r); err != nil {
|
2020-04-06 19:03:42 +00:00
|
|
|
return nil, err
|
|
|
|
}
|
|
|
|
|
|
|
|
// Copy the bytes so we could empty the buffer for reuse.
|
2020-10-28 17:19:22 +00:00
|
|
|
cpy := make([]byte, state.buf.Len())
|
|
|
|
copy(cpy, state.buf.Bytes())
|
2020-08-20 21:15:52 +00:00
|
|
|
|
|
|
|
// If the buffer's capacity is over the limit, then re-allocate a new one.
|
2020-10-28 17:19:22 +00:00
|
|
|
if state.buf.Cap() > MaxCapUntilReset {
|
|
|
|
state.buf = bytes.Buffer{}
|
|
|
|
state.buf.Grow(CopyBufferSize)
|
2020-08-20 21:15:52 +00:00
|
|
|
}
|
2020-04-06 19:03:42 +00:00
|
|
|
|
|
|
|
return cpy, nil
|
|
|
|
}
|