gws

package module
v1.4.5 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Apr 11, 2023 License: MIT Imports: 21 Imported by: 64

README

gws

event-driven go websocket server

Build Status MIT licensed Go Version codecov Go Report Card

Highlight
  • No dependency
  • No channel, no additional resident concurrent goroutine
  • Asynchronous non-blocking read and write support
  • High IOPS and low latency
  • Fully passes the WebSocket autobahn-testsuite
Install
go get -v github.com/lxzan/gws@latest
Interface
type Event interface {
	OnOpen(socket *Conn)
	OnError(socket *Conn, err error)
	OnClose(socket *Conn, code uint16, reason []byte)
	OnPing(socket *Conn, payload []byte)
	OnPong(socket *Conn, payload []byte)
	OnMessage(socket *Conn, message *Message)
}
Examples
Server
package main

import (
	"fmt"
	"github.com/lxzan/gws"
	"net/http"
)

func main() {
	var upgrader = gws.NewUpgrader(new(WebSocket), &gws.ServerOption{
		CompressEnabled:     true,
		CheckUtf8Enabled:    true,
		ReadMaxPayloadSize:  32 * 1024 * 1024,
		WriteMaxPayloadSize: 32 * 1024 * 1024,
		ReadAsyncEnabled:    true,
		ReadBufferSize:      4 * 1024,
		WriteBufferSize:     4 * 1024,
	})

	http.HandleFunc("/connect", func(writer http.ResponseWriter, request *http.Request) {
		socket, err := upgrader.Accept(writer, request)
		if err != nil {
			return
		}
		socket.Listen()
	})

	_ = http.ListenAndServe(":3000", nil)
}

type WebSocket struct{}

func (c *WebSocket) OnClose(socket *gws.Conn, code uint16, reason []byte) {
	fmt.Printf("onclose: code=%d, payload=%s\n", code, string(reason))
}

func (c *WebSocket) OnError(socket *gws.Conn, err error) {
	fmt.Printf("onerror: err=%s\n", err.Error())
}

func (c *WebSocket) OnOpen(socket *gws.Conn) {
	println("connected")
}

func (c *WebSocket) OnPing(socket *gws.Conn, payload []byte) {
	fmt.Printf("onping: payload=%s\n", string(payload))
	socket.WritePong(payload)
}

func (c *WebSocket) OnPong(socket *gws.Conn, payload []byte) {}

func (c *WebSocket) OnMessage(socket *gws.Conn, message *gws.Message) {
	defer message.Close()
	socket.WriteMessage(message.Opcode, message.Data.Bytes())
}
Client
package main

import (
	"fmt"
	"github.com/lxzan/gws"
	"log"
)

func main() {
	socket, _, err := gws.NewClient(new(WebSocket), &gws.ClientOption{
		Addr: "ws://127.0.0.1:3000/connect",
	})
	if err != nil {
		log.Printf(err.Error())
		return
	}
	socket.Listen()
}

type WebSocket struct {
	gws.BuiltinEventHandler
}

func (c *WebSocket) OnMessage(socket *gws.Conn, message *gws.Message) {
	fmt.Printf("recv: %s\n", message.Data.String())
}
TLS
package main

import (
	"github.com/gin-gonic/gin"
	"github.com/lxzan/gws"
)

func main() {
	app := gin.New()
	handler := new(WebSocket)
	upgrader := gws.NewUpgrader(handler, nil)
	app.GET("/connect", func(ctx *gin.Context) {
		socket, err := upgrader.Accept(ctx.Writer, ctx.Request)
		if err != nil {
			return
		}
		upgrader.Listen(socket)
	})
	cert := "server.crt"
	key := "server.key"
	if err := app.RunTLS(":8443", cert, key); err != nil {
		panic(err)
	}
}
Autobahn Test
cd examples/testsuite
mkdir reports
docker run -it --rm \
    -v ${PWD}/config:/config \
    -v ${PWD}/reports:/reports \
    crossbario/autobahn-testsuite \
    wstest -m fuzzingclient -s /config/fuzzingclient.json
Benchmark
  • Machine: Ubuntu 20.04LTS VM (4C8T)
Max IOPS
tcpkali -c 1000 --connect-rate 500 -r ${message_num} -T 300s -f assets/1K.txt --ws 127.0.0.1:${port}/connect

rps

Latency
tcpkali -c 1000 --connect-rate 500 -r 100 -T 300s -f assets/1K.txt --ws 127.0.0.1:${port}/connect

gws-c1000-m100

gorilla-c1000-m100

CPU
  PID USER      PR  NI    VIRT    RES    SHR S  %CPU  %MEM     TIME+ COMMAND
26054 caster    20   0  720164  39320   7340 S 246.5   1.0  48:34.38 gorilla-linux-a
26059 caster    20   0  720852  53624   7196 S 179.4   1.3  48:39.85 gws-linux-amd64
Acknowledgments

The following project had particular influence on gws's design.

Documentation

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type BuiltinEventHandler added in v1.3.0

type BuiltinEventHandler struct{}

func (BuiltinEventHandler) OnClose added in v1.3.0

func (b BuiltinEventHandler) OnClose(socket *Conn, code uint16, reason []byte)

func (BuiltinEventHandler) OnError added in v1.3.0

func (b BuiltinEventHandler) OnError(socket *Conn, err error)

func (BuiltinEventHandler) OnMessage added in v1.3.0

func (b BuiltinEventHandler) OnMessage(socket *Conn, message *Message)

func (BuiltinEventHandler) OnOpen added in v1.3.0

func (b BuiltinEventHandler) OnOpen(socket *Conn)

func (BuiltinEventHandler) OnPing added in v1.3.0

func (b BuiltinEventHandler) OnPing(socket *Conn, payload []byte)

func (BuiltinEventHandler) OnPong added in v1.3.0

func (b BuiltinEventHandler) OnPong(socket *Conn, payload []byte)

type ClientOption added in v1.4.2

type ClientOption struct {
	// 写缓冲区的大小, v1.4.5版本此参数被废弃
	// Deprecated: Size of the write buffer, v1.4.5 version of this parameter is deprecated
	WriteBufferSize     int
	ReadAsyncEnabled    bool
	ReadAsyncGoLimit    int
	ReadAsyncCap        int
	ReadMaxPayloadSize  int
	ReadBufferSize      int
	WriteAsyncCap       int
	WriteMaxPayloadSize int
	CompressEnabled     bool
	CompressLevel       int
	CompressThreshold   int
	CheckUtf8Enabled    bool

	// 连接地址, 例如 wss://example.com/connect
	// service address, eg: wss://example.com/connect
	Addr string
	// 额外的请求头
	// extra request header
	RequestHeader http.Header
	// dial timeout
	// 连接超时时间
	DialTimeout time.Duration
	// TLS设置
	// tls config
	TlsConfig *tls.Config
}

type ConcurrentMap added in v1.2.5

type ConcurrentMap struct {
	// contains filtered or unexported fields
}

ConcurrentMap used to store websocket connections in the IM server 用来存储IM等服务的连接

func NewConcurrentMap added in v1.2.5

func NewConcurrentMap(segments uint64) *ConcurrentMap

func (*ConcurrentMap) Delete added in v1.2.5

func (c *ConcurrentMap) Delete(key interface{})

func (*ConcurrentMap) Len added in v1.2.5

func (c *ConcurrentMap) Len() int

func (*ConcurrentMap) Load added in v1.2.5

func (c *ConcurrentMap) Load(key interface{}) (value interface{}, exist bool)

func (*ConcurrentMap) Range added in v1.2.5

func (c *ConcurrentMap) Range(f func(key interface{}, value interface{}) bool)

Range calls f sequentially for each key and value present in the map. If f returns false, range stops the iteration.

func (*ConcurrentMap) Store added in v1.2.5

func (c *ConcurrentMap) Store(key interface{}, value interface{})

type Config added in v1.2.0

type Config struct {
	// 是否开启异步读, 开启的话会并行调用OnMessage
	// Whether to enable asynchronous reading, if enabled OnMessage will be called in parallel
	ReadAsyncEnabled bool

	// 异步读的最大并行协程数量
	// Maximum number of parallel concurrent processes for asynchronous reads
	ReadAsyncGoLimit int

	// 异步读的容量限制, 容量溢出将会返回错误
	// Capacity limit for asynchronous reads, overflow will return an error
	ReadAsyncCap int

	// 最大读取的消息内容长度
	// Maximum read message content length
	ReadMaxPayloadSize int

	// 读缓冲区的大小
	// Size of the read buffer
	ReadBufferSize int

	// 异步写的容量限制, 容量溢出将会返回错误
	// Capacity limit for asynchronous writes, overflow will return an error
	WriteAsyncCap int

	// 最大写入的消息内容长度
	// Maximum length of written message content
	WriteMaxPayloadSize int

	// 写缓冲区的大小, v1.4.5版本此参数被废弃
	// Deprecated: Size of the write buffer, v1.4.5 version of this parameter is deprecated
	WriteBufferSize int

	// 是否开启数据压缩
	// Whether to turn on data compression
	CompressEnabled bool

	// 压缩级别
	// Compress level
	CompressLevel int

	// 压缩阈值, 低于阈值的消息不会被压缩
	// Compression threshold, messages below the threshold will not be compressed
	CompressThreshold int

	// 是否检查文本utf8编码, 关闭性能会好点
	// Whether to check the text utf8 encoding, turn off the performance will be better
	CheckUtf8Enabled bool
}

type Conn

type Conn struct {
	// store session information
	SessionStorage SessionStorage
	// contains filtered or unexported fields
}

func NewClient added in v1.4.2

func NewClient(handler Event, option *ClientOption) (client *Conn, resp *http.Response, e error)

NewClient 创建WebSocket客户端

func (*Conn) Listen added in v1.1.2

func (c *Conn) Listen()

Listen listening to websocket messages through a dead loop 监听websocket消息

func (*Conn) LocalAddr added in v1.0.1

func (c *Conn) LocalAddr() net.Addr

func (*Conn) NetConn added in v1.2.10

func (c *Conn) NetConn() net.Conn

NetConn get tcp/tls/... conn

func (*Conn) RemoteAddr added in v1.0.1

func (c *Conn) RemoteAddr() net.Addr

func (*Conn) SetDeadline

func (c *Conn) SetDeadline(t time.Time) error

SetDeadline sets deadline

func (*Conn) SetReadDeadline added in v1.1.2

func (c *Conn) SetReadDeadline(t time.Time) error

SetReadDeadline sets read deadline

func (*Conn) SetWriteDeadline added in v1.1.2

func (c *Conn) SetWriteDeadline(t time.Time) error

SetWriteDeadline sets write deadline

func (*Conn) WriteAsync added in v1.3.0

func (c *Conn) WriteAsync(opcode Opcode, payload []byte) error

WriteAsync 异步非阻塞地写入消息 Write messages asynchronously and non-blockingly

func (*Conn) WriteClose

func (c *Conn) WriteClose(code uint16, reason []byte)

WriteClose proactively close the connection code: https://developer.mozilla.org/zh-CN/docs/Web/API/CloseEvent#status_codes 通过emitError发送关闭帧, 将连接状态置为关闭, 用于服务端主动断开连接 没有特殊原因的话, 建议code=0, reason=nil

func (*Conn) WriteMessage added in v1.1.0

func (c *Conn) WriteMessage(opcode Opcode, payload []byte) error

WriteMessage 发送消息 如果是客户端, payload内容会被改变 writes message

func (*Conn) WritePing

func (c *Conn) WritePing(payload []byte) error

WritePing write ping frame

func (*Conn) WritePong

func (c *Conn) WritePong(payload []byte) error

WritePong write pong frame

func (*Conn) WriteString added in v1.2.10

func (c *Conn) WriteString(s string) error

WriteString write text frame force convert string to []byte

type Event added in v1.1.2

type Event interface {
	// 建立连接事件
	OnOpen(socket *Conn)

	// 错误事件
	// IO错误, 协议错误, 压缩解压错误...
	OnError(socket *Conn, err error)

	// 关闭事件
	// 另一端发送了关闭帧
	OnClose(socket *Conn, code uint16, reason []byte)

	// 心跳探测事件
	OnPing(socket *Conn, payload []byte)

	// 心跳响应事件
	OnPong(socket *Conn, payload []byte)

	// 消息事件
	// 如果开启了AsyncReadEnabled, 可以在一个连接里面并行处理多个请求
	OnMessage(socket *Conn, message *Message)
}

WebSocket Event one of onclose and onerror will be called once during the connection's lifetime. 在连接的生命周期中,onclose和onerror中的一个有且只有一次被调用.

type Message

type Message struct {
	Opcode Opcode        // 帧状态码
	Data   *bytes.Buffer // 数据缓冲
}

func (*Message) Bytes

func (c *Message) Bytes() []byte

func (*Message) Close

func (c *Message) Close()

Close recycle buffer

func (*Message) Read added in v1.1.0

func (c *Message) Read(p []byte) (n int, err error)

type Opcode

type Opcode uint8
const (
	OpcodeContinuation    Opcode = 0x0
	OpcodeText            Opcode = 0x1
	OpcodeBinary          Opcode = 0x2
	OpcodeCloseConnection Opcode = 0x8
	OpcodePing            Opcode = 0x9
	OpcodePong            Opcode = 0xA
)

func (Opcode) IsDataFrame added in v1.1.2

func (c Opcode) IsDataFrame() bool

type ServerOption added in v1.4.0

type ServerOption struct {
	// 写缓冲区的大小, v1.4.5版本此参数被废弃
	// Deprecated: Size of the write buffer, v1.4.5 version of this parameter is deprecated
	WriteBufferSize     int
	ReadAsyncEnabled    bool
	ReadAsyncGoLimit    int
	ReadAsyncCap        int
	ReadMaxPayloadSize  int
	ReadBufferSize      int
	WriteAsyncCap       int
	WriteMaxPayloadSize int
	CompressEnabled     bool
	CompressLevel       int
	CompressThreshold   int
	CheckUtf8Enabled    bool

	// WebSocket子协议, 一般不需要设置
	// WebSocket subprotocol, usually no need to set
	Subprotocols []string

	// 连接握手时添加的额外的响应头, 如果客户端不支持就不要传
	// https://www.rfc-editor.org/rfc/rfc6455.html#section-1.3
	// attention: client may not support custom response header, use nil instead
	ResponseHeader http.Header

	// 检查请求来源
	// Check the origin of the request
	CheckOrigin func(r *http.Request, session SessionStorage) bool
}

type SessionStorage added in v1.2.3

type SessionStorage interface {
	Load(key string) (value interface{}, exist bool)
	Delete(key string)
	Store(key string, value interface{})
	Range(f func(key string, value interface{}) bool)
}

SessionStorage because sync.Map is not easy to debug, so I implemented my own map. if you don't like it, use sync.Map instead.

type Upgrader

type Upgrader struct {
	// contains filtered or unexported fields
}

func NewUpgrader added in v1.2.11

func NewUpgrader(eventHandler Event, option *ServerOption) *Upgrader

func (*Upgrader) Accept added in v1.2.11

func (c *Upgrader) Accept(w http.ResponseWriter, r *http.Request) (*Conn, error)

Accept http upgrade to websocket protocol

Directories

Path Synopsis
examples
chatroom command
client command
server command
testsuite command

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL