msock

package module
v1.0.25 Latest Latest
Warning

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

Go to latest
Published: Aug 12, 2026 License: Apache-2.0 Imports: 18 Imported by: 0

README

msock

统一封装 TCP / WebSocket / KCP 的网络通信库,提供一致的服务端与客户端 API、消息路由、中间件链和多种内置编解码器。

安装

go get github.com/quietking0312/component/msock

快速开始

服务端
router := msock.NewRouter()
router.Register(1, func(conn msock.Conn, msg msock.Message) {
    fmt.Println("收到消息:", string(msg.Data()))
    conn.Send(msock.NewMessage(2, []byte("pong")))
})

server, err := msock.NewServer(
    msock.WithAddress(":8080"),
    msock.WithConnType(msock.ConnTypeTCP),
)
if err != nil {
    log.Fatal(err)
}
server.SetRouter(router)
log.Fatal(server.Run())
客户端
router := msock.NewRouter()
router.Register(2, func(conn msock.Conn, msg msock.Message) {
    fmt.Println("收到回复:", string(msg.Data()))
})

client := msock.NewClient(
    msock.WithConnType(msock.ConnTypeTCP),
)
client.SetRouter(router)

if err := client.Connect(":8080"); err != nil {
    log.Fatal(err)
}
defer client.Close()
client.Send(msock.NewMessage(1, []byte("ping")))

支持的协议

常量 协议 说明
ConnTypeTCP TCP 原生 TCP 长连接
ConnTypeWebSocket WebSocket gorilla/websocket
ConnTypeGWS WebSocket lxzan/gws(高性能)
ConnTypeKCP KCP 基于 UDP 的可靠传输

协议切换只需修改 WithConnType,其余代码不变。WebSocket 服务端默认挂载 /ws 路径。

消息

Message 接口包含两个字段:RouteID() uint32 用于路由,Data() []byte 为消息体。

msg := msock.NewMessage(1, []byte("hello"))
msg.RouteID() // 1
msg.Data()    // []byte("hello")

高并发场景可使用对象池减少 GC 压力:

msg := msock.AcquireMessage()
msg.SetData([]byte("hello"))
// ...
msock.ReleaseMessage(msg)

编解码器

内置编解码器

SimpleCodec(默认)— 二进制协议

[4字节总长度 BE] [4字节RouteID BE] [Data]
msock.NewSimpleCodec()           // 默认 64KB 限制
msock.NewSimpleCodec(128 * 1024) // 自定义最大包大小

TLVCodec — Type 即 RouteID,适合消息类型 ≤ 255 的场景

[1字节Type] [2字节Length BE] [Value]
msock.NewTLVCodec()

LineCodec — 文本协议,RouteID 固定为 0

[文本内容]\n
msock.NewLineCodec()
自定义编解码器

实现 Codec 接口即可:

type Codec interface {
    Encode(msg Message) ([]byte, error)
    Decode(data []byte) (Message, int, error) // 返回消息和已消费字节数
    MaxPacketSize() int
}

路由

router := msock.NewRouter()

// 注册处理器
router.Register(1, handleLogin)
router.Register(2, handleLogout)

// 批量注册
router.RegisterMultiple(map[uint32]msock.Handler{
    3: handleA,
    4: handleB,
})

// 移除
router.Remove(1)

// 自定义未匹配处理器
router.SetNotFoundHandler(func(conn msock.Conn, msg msock.Message) {
    fmt.Println("未知路由:", msg.RouteID())
})

中间件

中间件签名为 func(next Handler) Handler,通过 router.Use() 全局注册,按注册顺序形成洋葱模型。

router.Use(msock.Recovery(logger), msock.Logging(logger))
内置中间件
// panic 恢复,记录错误日志后继续运行
msock.Recovery(logger)

// 请求日志,记录 conn/route/size
msock.Logging(logger)

// 认证,返回 false 则中断后续处理
msock.Auth(func(conn msock.Conn) bool {
    _, ok := conn.GetValue("uid")
    return ok
}, logger)

// 消息校验
msock.Validate(func(msg msock.Message) error {
    if len(msg.Data()) == 0 {
        return errors.New("empty data")
    }
    return nil
}, logger)

// 限流(每连接计数器)
msock.RateLimit(100, logger)

// 超时(连接 context 取消时触发)
msock.Timeout(func() { fmt.Println("超时") }, logger)
自定义中间件
func myMiddleware(next msock.Handler) msock.Handler {
    return func(conn msock.Conn, msg msock.Message) {
        // 前置逻辑
        next(conn, msg)
        // 后置逻辑
    }
}

router.Use(myMiddleware)

连接事件

// 服务端
server.OnConnect(func(conn msock.Conn) {
    fmt.Println("连接建立:", conn.ID(), conn.Type())
})
server.OnDisconnect(func(conn msock.Conn) {
    fmt.Println("连接断开:", conn.ID())
})
server.OnError(func(conn msock.Conn, err error) {
    fmt.Println("连接错误:", err)
})

// 客户端
client.OnConnect(func(conn msock.Conn) { ... })
client.OnDisconnect(func(conn msock.Conn) { ... })
client.OnError(func(err error) { ... })

无论是对端主动断开、网络超时还是服务端调用 conn.Close()OnDisconnect 都会触发。

连接对象

conn.ID()           // 唯一标识(UUID)
conn.Type()         // ConnTypeTCP / ConnTypeWebSocket / ConnTypeGWS / ConnTypeKCP
conn.RemoteAddr()   // 远端地址
conn.Send(msg)      // 发送消息(异步入队)
conn.SendBytes(b)   // 发送原始字节
conn.Close()        // 关闭连接
conn.IsClosed()     // 是否已关闭
conn.Context()      // 连接的 context,断开时 Done()

// 跨 handler 传递数据
conn.SetValue("uid", 12345)
uid, ok := conn.GetValue("uid")

连接管理

mgr := server.GetConnManager()
mgr.Count()          // 当前连接数
mgr.Get("conn-id")   // 获取单个连接
mgr.GetAll()         // 获取所有连接

server.SendTo("conn-id", msg) // 发送给指定连接
server.Broadcast(msg)         // 广播(编码一次,批量发送)
server.ConnCount()            // 当前连接数

服务端配置

选项 默认值 说明
WithAddress :8080 监听地址
WithConnType ConnTypeTCP 协议类型
WithCodec SimpleCodec 编解码器
WithLogger 空日志 日志实现
WithMaxConnections 10000 最大连接数
WithReadBufferSize 4096 读缓冲区字节数
WithWriteBufferSize 4096 写缓冲区字节数
WithReadTimeout 60s 读超时
WithWriteTimeout 10s 写超时
WithHeartbeat 30s / 90s 心跳间隔 / 超时
WithGWSConfig nil gws 服务端详细配置
GWS 配置

ConnTypeConnTypeGWS 时,gws.ServerOption 可以通过配置文件决定。

msock.GWSConfig 已内置 json / yaml / mapstructure 标签,你可以使用任意配置库(如 encoding/jsongopkg.in/yaml.v3、viper 等)反序列化到该结构体,再通过 WithGWSConfig 传入。

配置文件示例 gws.yaml

gws:
  read_buffer_size: 8192
  read_max_payload_size: 16384
  write_max_payload_size: 16384
  parallel_enabled: true
  parallel_golimit: 0
  check_utf8_enabled: false
  handshake_timeout: 10s
  sub_protocols:
    - chat
  response_header:
    X-Custom:
      - value
  permessage_deflate:
    enabled: false

使用示例:

gwsCfg := &msock.GWSConfig{}
// 使用你自己的配置工具加载到 gwsCfg,例如:
// yaml.Unmarshal(data, gwsCfg)

server, err := msock.NewServer(
    msock.WithAddress(":8080"),
    msock.WithConnType(msock.ConnTypeGWS),
    msock.WithGWSConfig(gwsCfg),
)

未显式配置的字段可先用 msock.DefaultGWSConfig() 作为基础,再反序列化覆盖。如果完全不传 WithGWSConfig,则保持与历史行为一致的默认值。

WebSocket 原生 ping/pong 帧也可以通过 GWSConfig 自定义:

gwsCfg.OnPing = func(payload []byte) []byte {
    // 处理对端发来的 ping payload,返回的值会作为 pong 帧回给对方
    return []byte("pong")
}
gwsCfg.OnPong = func(payload []byte) {
    // 处理对端发来的 pong payload
}

注意:OnPing / OnPong 控制的是 WebSocket 协议层的 ping/pong 帧;msock 自身的应用层心跳(HeartbeatPingID / HeartbeatPongID)不受它们影响。

日志

实现 Logger 接口对接任意日志库:

type Logger interface {
    Debug(msg string, args ...any)
    Info(msg string, args ...any)
    Warn(msg string, args ...any)
    Error(msg string, args ...any)
}

内置提供基于标准库 logStdLogger

msock.WithLogger(msock.NewStdLogger())

Documentation

Index

Constants

This section is empty.

Variables

View Source
var (
	// ErrConnClosed 连接已关闭
	ErrConnClosed = errors.New("connection closed")

	// ErrServerClosed 服务器已关闭
	ErrServerClosed = errors.New("server closed")

	// ErrInvalidMessage 无效消息
	ErrInvalidMessage = errors.New("invalid message")

	// ErrCodecNotSet 编解码器未设置
	ErrCodecNotSet = errors.New("codec not set")

	// ErrRouterNotSet 路由器未设置
	ErrRouterNotSet = errors.New("router not set")

	// ErrMaxConnections 达到最大连接数
	ErrMaxConnections = errors.New("max connections reached")

	// ErrListenFailed 监听失败
	ErrListenFailed = errors.New("listen failed")

	// ErrAcceptFailed 接受连接失败
	ErrAcceptFailed = errors.New("accept connection failed")

	// ErrTimeout 超时
	ErrTimeout = errors.New("timeout")

	// ErrPacketTooLarge 数据包太大
	ErrPacketTooLarge = errors.New("packet too large")

	// ErrInvalidPacket 无效数据包
	ErrInvalidPacket = errors.New("invalid packet")

	// ErrUnsupportedProtocol 不支持的协议
	ErrUnsupportedProtocol = errors.New("unsupported protocol")

	// ErrSendChannelFull 发送通道已满
	ErrSendChannelFull = errors.New("send channel full")

	// ErrNoAvailableConn 连接池中无可用连接
	ErrNoAvailableConn = errors.New("no available connection in pool")
)

Functions

func ReleaseMessage

func ReleaseMessage(msg *DefaultMessage)

ReleaseMessage 将消息对象放回池中

Types

type Client

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

Client 网络客户端(连接池 + 断线重连)

func NewClient

func NewClient(opts ...ServerOption) *Client

NewClient 创建客户端

func (*Client) AvailableConns added in v1.0.19

func (c *Client) AvailableConns() int

AvailableConns 返回当前可用(已连接)的连接数

func (*Client) Close

func (c *Client) Close() error

Close 关闭客户端及所有连接,等待后台 goroutine 退出

func (*Client) Conn

func (c *Client) Conn() Conn

Conn 返回第一个可用连接(兼容旧接口)

func (*Client) Connect

func (c *Client) Connect(addr string) error

Connect 连接到服务器,建立连接池中所有连接

func (*Client) HandleHeartbeat added in v1.0.9

func (c *Client) HandleHeartbeat(conn Conn, msg Message)

HandleHeartbeat 客户端收到 Pong 时调用,更新心跳时间戳 需要在客户端 Router 中注册 HeartbeatPongID 路由时调用此函数

func (*Client) IsClosed

func (c *Client) IsClosed() bool

IsClosed 检查客户端是否已关闭

func (*Client) OnConnect

func (c *Client) OnConnect(fn func(Conn))

OnConnect 设置连接建立回调

func (*Client) OnDisconnect

func (c *Client) OnDisconnect(fn func(Conn))

OnDisconnect 设置连接断开回调

func (*Client) OnError

func (c *Client) OnError(fn func(error))

OnError 设置错误回调

func (*Client) PoolSize added in v1.0.19

func (c *Client) PoolSize() int

PoolSize 返回连接池大小

func (*Client) RegisterHeartbeat added in v1.0.9

func (c *Client) RegisterHeartbeat(router *Router)

RegisterHeartbeat 向 Router 注册客户端心跳 Pong 处理器

func (*Client) Send

func (c *Client) Send(msg Message) error

Send 轮询连接池发送消息

func (*Client) SendBytes

func (c *Client) SendBytes(data []byte) error

SendBytes 轮询连接池发送原始字节

func (*Client) SetCodec

func (c *Client) SetCodec(codec Codec)

SetCodec 设置编解码器

func (*Client) SetRouter

func (c *Client) SetRouter(router *Router)

SetRouter 设置路由器

func (*Client) Stats added in v1.0.19

func (c *Client) Stats() ClientStats

Stats 返回客户端统计快照

type ClientStats added in v1.0.19

type ClientStats struct {
	// PoolSize 连接池槽位总数
	PoolSize int
	// AvailableConns 当前可用(已连接)的连接数
	AvailableConns int
	// TotalReconnects 累计重连次数
	TotalReconnects int64
	// TotalMessages 累计处理的消息数
	TotalMessages int64
	// TotalRecvBytes 累计接收字节数
	TotalRecvBytes int64
	// TotalSendBytes 累计发送字节数
	TotalSendBytes int64
}

ClientStats 客户端统计快照

type Codec

type Codec interface {
	// Encode 编码消息为字节流
	Encode(msg Message) ([]byte, error)
	// HeaderSize 返回固定 header 字节数
	HeaderSize() int
	// DecodeHeader 解析 header,返回 routeID 和 body 长度
	DecodeHeader(header []byte) (routeID uint32, bodyLen int, err error)
	// DecodeBody 将 body 字节解析为消息
	DecodeBody(routeID uint32, body []byte) (Message, error)
	// MaxPacketSize 返回允许的最大包大小(header + body)
	MaxPacketSize() int
}

Codec 编解码器接口,采用两阶段解码:先解 header 得到 body 长度,再读 body。

type Conn

type Conn interface {
	// ID 返回连接唯一标识
	ID() string
	// Type 返回连接类型
	Type() ConnType
	// LocalAddr 返回本地地址
	LocalAddr() net.Addr
	// RemoteAddr 返回远程地址
	RemoteAddr() net.Addr
	// Send 发送消息
	Send(msg Message) error
	// SendBytes 发送原始字节数据
	SendBytes(data []byte) error
	// Close 关闭连接
	Close() error
	// IsClosed 检查连接是否已关闭
	IsClosed() bool
	// SetReadDeadline 设置读取超时
	SetReadDeadline(t time.Time) error
	// SetWriteDeadline 设置写入超时
	SetWriteDeadline(t time.Time) error
	// Context 返回连接的上下文
	Context() context.Context
	// SetValue 存储键值对到连接上下文
	SetValue(key, value interface{})
	// GetValue 从连接上下文获取值
	GetValue(key interface{}) (interface{}, bool)
}

Conn 连接接口,统一封装 TCP/WebSocket/KCP

type ConnManager

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

ConnManager 分片锁连接管理器 32 个分片,按连接 ID 首字节哈希,锁竞争降低为原来的 1/32

func NewConnManager

func NewConnManager(maxConn int) *ConnManager

NewConnManager 创建连接管理器

func (*ConnManager) Add

func (m *ConnManager) Add(conn Conn) bool

Add 添加连接,超过最大连接数返回 false

func (*ConnManager) BroadcastBytes

func (m *ConnManager) BroadcastBytes(data []byte) int

BroadcastBytes 广播原始字节到所有连接,返回成功入队的连接数

func (*ConnManager) BroadcastBytesFilter added in v1.0.25

func (m *ConnManager) BroadcastBytesFilter(data []byte, filter func(Conn) bool) int

BroadcastBytesFilter 广播原始字节到满足条件的连接,filter 为 nil 则广播到所有连接,返回成功入队的连接数

func (*ConnManager) BroadcastBytesTo added in v1.0.23

func (m *ConnManager) BroadcastBytesTo(data []byte, ids []string) int

BroadcastBytesTo 广播原始字节到指定 ID 的连接,返回成功入队的连接数

func (*ConnManager) CloseAll

func (m *ConnManager) CloseAll()

CloseAll 关闭所有连接

func (*ConnManager) Count

func (m *ConnManager) Count() int

Count 返回当前连接数

func (*ConnManager) Get

func (m *ConnManager) Get(connID string) (Conn, bool)

Get 获取连接

func (*ConnManager) GetAll

func (m *ConnManager) GetAll() []Conn

GetAll 获取所有连接快照

func (*ConnManager) Remove

func (m *ConnManager) Remove(connID string)

Remove 移除连接

type ConnType

type ConnType string

ConnType 连接类型

const (
	ConnTypeTCP       ConnType = "tcp"
	ConnTypeWebSocket ConnType = "websocket"
	ConnTypeKCP       ConnType = "kcp"
	ConnTypeGWS       ConnType = "gws"
)

type DefaultMessage

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

DefaultMessage 默认消息实现

func AcquireMessage

func AcquireMessage() *DefaultMessage

AcquireMessage 从池中获取消息对象

func NewMessage

func NewMessage(routeID uint32, data []byte) *DefaultMessage

NewMessage 创建新消息

func (*DefaultMessage) Data

func (m *DefaultMessage) Data() []byte

Data 返回消息数据

func (*DefaultMessage) RouteID

func (m *DefaultMessage) RouteID() uint32

RouteID 返回路由ID

func (*DefaultMessage) SetData

func (m *DefaultMessage) SetData(data []byte)

SetData 设置消息数据

type GWSConfig added in v1.0.14

type GWSConfig struct {
	// ReadBufferSize 读取缓冲区大小。
	ReadBufferSize int `mapstructure:"read_buffer_size" json:"read_buffer_size" yaml:"read_buffer_size"`
	// ReadMaxPayloadSize 读取最大负载大小。
	ReadMaxPayloadSize int `mapstructure:"read_max_payload_size" json:"read_max_payload_size" yaml:"read_max_payload_size"`
	// WriteBufferSize 写入缓冲区大小(gws 已废弃该参数,建议留空)。
	WriteBufferSize int `mapstructure:"write_buffer_size" json:"write_buffer_size" yaml:"write_buffer_size"`
	// WriteMaxPayloadSize 写入最大负载大小。
	WriteMaxPayloadSize int `mapstructure:"write_max_payload_size" json:"write_max_payload_size" yaml:"write_max_payload_size"`
	// ParallelEnabled 是否启用并行处理。
	ParallelEnabled bool `mapstructure:"parallel_enabled" json:"parallel_enabled" yaml:"parallel_enabled"`
	// ParallelGolimit 并行协程限制,<=0 表示不限制。
	ParallelGolimit int `mapstructure:"parallel_golimit" json:"parallel_golimit" yaml:"parallel_golimit"`
	// CheckUtf8Enabled 是否启用 UTF-8 检查。
	CheckUtf8Enabled bool `mapstructure:"check_utf8_enabled" json:"check_utf8_enabled" yaml:"check_utf8_enabled"`
	// HandshakeTimeout 握手超时时间。
	HandshakeTimeout time.Duration `mapstructure:"handshake_timeout" json:"handshake_timeout" yaml:"handshake_timeout"`
	// SubProtocols WebSocket 子协议列表。
	SubProtocols []string `mapstructure:"sub_protocols" json:"sub_protocols" yaml:"sub_protocols"`
	// ResponseHeader 握手时附加的响应头。
	ResponseHeader http.Header `mapstructure:"response_header" json:"response_header" yaml:"response_header"`
	// PermessageDeflate 压缩扩展配置。
	PermessageDeflate GWSPermessageDeflateConfig `mapstructure:"permessage_deflate" json:"permessage_deflate" yaml:"permessage_deflate"`

	// Recovery 自定义 panic 恢复函数,为 nil 时使用 gws.Recovery。
	Recovery func(logger gws.Logger) `mapstructure:"-" json:"-" yaml:"-"`
	// Authorize 自定义鉴权函数,为 nil 时允许所有请求。
	Authorize func(r *http.Request, session gws.SessionStorage) bool `mapstructure:"-" json:"-" yaml:"-"`
	// NewSession 自定义 SessionStorage 工厂,为 nil 时使用 gws 默认实现。
	NewSession func() gws.SessionStorage `mapstructure:"-" json:"-" yaml:"-"`

	// OnPing 收到 WebSocket ping 帧时的回调。
	// 返回的 payload 会作为 pong 帧回复给对端;为 nil 时使用默认空 pong。
	OnPing func(payload []byte) []byte `mapstructure:"-" json:"-" yaml:"-"`
	// OnPong 收到 WebSocket pong 帧时的回调。
	OnPong func(payload []byte) `mapstructure:"-" json:"-" yaml:"-"`
}

GWSConfig 用于配置 gws 服务端选项。 零值或 nil 表示使用 msock 的默认行为(与 ServerConfig 中的缓冲区大小保持一致)。

func DefaultGWSConfig added in v1.0.14

func DefaultGWSConfig() *GWSConfig

DefaultGWSConfig 返回与历史默认行为一致的 GWS 配置。

type GWSPermessageDeflateConfig added in v1.0.14

type GWSPermessageDeflateConfig struct {
	// Enabled 是否开启压缩。
	Enabled bool `mapstructure:"enabled" json:"enabled" yaml:"enabled"`
	// Level 压缩级别。
	Level int `mapstructure:"level" json:"level" yaml:"level"`
	// Threshold 压缩阈值,长度小于阈值的消息不会被压缩。
	Threshold int `mapstructure:"threshold" json:"threshold" yaml:"threshold"`
	// PoolSize 压缩器内存池大小。
	PoolSize int `mapstructure:"pool_size" json:"pool_size" yaml:"pool_size"`
	// ServerContextTakeover 服务端上下文接管。
	ServerContextTakeover bool `mapstructure:"server_context_takeover" json:"server_context_takeover" yaml:"server_context_takeover"`
	// ClientContextTakeover 客户端上下文接管。
	ClientContextTakeover bool `mapstructure:"client_context_takeover" json:"client_context_takeover" yaml:"client_context_takeover"`
	// ServerMaxWindowBits 服务端滑动窗口指数(8~15)。
	ServerMaxWindowBits int `mapstructure:"server_max_window_bits" json:"server_max_window_bits" yaml:"server_max_window_bits"`
	// ClientMaxWindowBits 客户端滑动窗口指数(8~15)。
	ClientMaxWindowBits int `mapstructure:"client_max_window_bits" json:"client_max_window_bits" yaml:"client_max_window_bits"`
}

GWSPermessageDeflateConfig 是 gws.PermessageDeflate 的可配置子集。

type Group

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

Group 路由组

func (*Group) Register

func (g *Group) Register(routeID uint32, handler Handler)

Register 在组内注册消息处理器 routeID: 组内的相对路由ID,最终路由ID = prefix << 16 | routeID

func (*Group) Use

func (g *Group) Use(mw ...Middleware)

Use 为路由组添加中间件

type Handler

type Handler func(conn Conn, msg Message)

Handler 消息处理器

type KCPConfig

type KCPConfig struct {
	// 发送窗口大小
	SendWindow int
	// 接收窗口大小
	RecvWindow int
	// 数据包最大传输单元
	Mtu int
	// 是否启用 FEC
	EnableFEC bool
	// FEC 数据分片数
	DataShards int
	// FEC 校验分片数
	ParityShards int
	// 是否启用加密
	EnableCrypt bool
	// 加密密钥(EnableCrypt 为 true 时必填)
	CryptKey string

	// 以下对应 kcp.UDPSession.SetNoDelay 的四个参数
	// NoDelay: 0=关闭,1=开启 nodelay 模式
	NoDelay int
	// Interval: 内部刷新时间间隔(毫秒)
	Interval int
	// Resend: 快速重传模式,0=关闭,2=推荐值
	Resend int
	// NC: 是否关闭流量控制,0=开启,1=关闭
	NC int
}

KCPConfig KCP配置

func DefaultKCPConfig

func DefaultKCPConfig() *KCPConfig

DefaultKCPConfig 返回默认KCP配置

type LineCodec

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

LineCodec 文本行编解码器,适用于文本协议(如 Telnet) 因行长度不固定,HeaderSize 返回 0,bufferedReader 走流式扫描路径。

func NewLineCodec

func NewLineCodec(maxLineLength ...int) *LineCodec

NewLineCodec 创建文本行编解码器

func (*LineCodec) DecodeBody added in v1.0.9

func (c *LineCodec) DecodeBody(routeID uint32, body []byte) (Message, error)

DecodeBody 不适用于 LineCodec

func (*LineCodec) DecodeHeader added in v1.0.9

func (c *LineCodec) DecodeHeader(header []byte) (uint32, int, error)

DecodeHeader 不适用于 LineCodec,始终返回错误

func (*LineCodec) Encode

func (c *LineCodec) Encode(msg Message) ([]byte, error)

Encode 编码文本行(末尾添加换行符)

func (*LineCodec) HeaderSize added in v1.0.9

func (c *LineCodec) HeaderSize() int

HeaderSize 返回 0,表示使用流式扫描而非固定 header

func (*LineCodec) MaxPacketSize

func (c *LineCodec) MaxPacketSize() int

MaxPacketSize 返回最大行长度

func (*LineCodec) ScanLine added in v1.0.9

func (c *LineCodec) ScanLine(data []byte) (Message, int, error)

ScanLine 扫描缓冲区中的一行,返回消息和已消费字节数;数据不足时返回 nil,0,nil

type Logger

type Logger interface {
	Debug(msg string, args ...any)
	Info(msg string, args ...any)
	Warn(msg string, args ...any)
	Error(msg string, args ...any)
}

Logger 日志接口,统一使用 slog 风格。

func NewStdLogger

func NewStdLogger() Logger

NewStdLogger 创建标准库日志记录器

type Message

type Message interface {
	// RouteID 返回路由ID,用于消息路由
	RouteID() uint32
	// Data 返回消息数据
	Data() []byte
	// SetData 设置消息数据
	SetData([]byte)
}

Message 消息接口

type Middleware

type Middleware func(next Handler) Handler

Middleware 中间件函数

func Auth

func Auth(authFunc func(conn Conn) bool, logger Logger) Middleware

Auth 认证中间件示例(需要配合连接上下文使用)

func Chain

func Chain(mws ...Middleware) Middleware

Chain 创建中间件链

func Logging

func Logging(logger Logger) Middleware

Logging 日志中间件

func RateLimit

func RateLimit(r float64, burst int, logger Logger) Middleware

RateLimit 限流中间件,基于令牌桶算法,按连接独立限流。 r: 每秒补充的令牌数(即每秒允许的请求速率) burst: 令牌桶容量(允许的瞬时突发量) 连接断开后自动清理限流器。

func Recovery

func Recovery(logger Logger) Middleware

Recovery 恢复中间件,捕获 panic

func Validate

func Validate(validateFunc func(msg Message) error, logger Logger) Middleware

Validate 消息校验中间件

type Router

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

Router 消息路由器

func NewRouter

func NewRouter() *Router

NewRouter 创建新的路由器

func (*Router) Get

func (r *Router) Get(routeID uint32) Handler

Get 获取指定路由ID的处理器

func (*Router) Group

func (r *Router) Group(prefix uint32, mws ...Middleware) *Group

Group 创建路由组,可以设置统一的前缀和中间件 prefix: 路由ID前缀,组内所有路由ID都会与该前缀组合

func (*Router) Handle

func (r *Router) Handle(conn Conn, msg Message)

Handle 处理消息

func (*Router) Register

func (r *Router) Register(routeID uint32, handler Handler)

Register 注册消息处理器

func (*Router) RegisterMultiple

func (r *Router) RegisterMultiple(handlers map[uint32]Handler)

RegisterMultiple 批量注册消息处理器

func (*Router) Remove

func (r *Router) Remove(routeID uint32)

Remove 移除指定路由ID的处理器

func (*Router) RouteCount

func (r *Router) RouteCount() int

RouteCount 返回注册的路由数量

func (*Router) SetNotFoundHandler

func (r *Router) SetNotFoundHandler(handler Handler)

SetNotFoundHandler 设置未找到路由时的处理器

func (*Router) Use

func (r *Router) Use(mw ...Middleware)

Use 添加中间件

type Server

type Server struct {
	BusinessWg sync.WaitGroup // 业务方在 handler 内异步 goroutine 的生命周期管理
	// contains filtered or unexported fields
}

Server 网络服务器

func NewServer

func NewServer(opts ...ServerOption) (*Server, error)

NewServer 创建服务器

func (*Server) Broadcast

func (s *Server) Broadcast(msg Message, ids ...string)

Broadcast 广播消息到指定 ID 的连接,ids 为空则广播到所有连接

func (*Server) BroadcastFilter added in v1.0.25

func (s *Server) BroadcastFilter(msg Message, filter func(Conn) bool)

BroadcastFilter 广播消息到满足条件的连接,filter 为 nil 则广播到所有连接

func (*Server) ConnCount

func (s *Server) ConnCount() int

ConnCount 返回当前连接数

func (*Server) GetConfig

func (s *Server) GetConfig() *ServerConfig

GetConfig 获取配置

func (*Server) GetConnManager

func (s *Server) GetConnManager() *ConnManager

GetConnManager 获取连接管理器

func (*Server) OnConnect

func (s *Server) OnConnect(fn func(Conn))

OnConnect 设置连接建立回调

func (*Server) OnDisconnect

func (s *Server) OnDisconnect(fn func(Conn))

OnDisconnect 设置连接断开回调

func (*Server) OnError

func (s *Server) OnError(fn func(Conn, error))

OnError 设置错误回调

func (*Server) RegisterHeartbeat added in v1.0.9

func (s *Server) RegisterHeartbeat(router *Router)

RegisterHeartbeat 向 Router 注册心跳处理器,需在 SetRouter 之后调用

func (*Server) RegisterWSRoute

func (s *Server) RegisterWSRoute(path string, upgrader *WebSocketUpgrader)

RegisterWSRoute 注册WebSocket处理路由(如果使用HTTP路由)

func (*Server) Run

func (s *Server) Run() error

Run 启动服务器

func (*Server) SendTo

func (s *Server) SendTo(connID string, msg Message) error

SendTo 发送消息到指定连接

func (*Server) SetCodec

func (s *Server) SetCodec(codec Codec)

SetCodec 设置编解码器

func (*Server) SetLogger

func (s *Server) SetLogger(logger Logger)

SetLogger 设置日志器

func (*Server) SetRouter

func (s *Server) SetRouter(router *Router)

SetRouter 设置路由器

func (*Server) Stats added in v1.0.19

func (s *Server) Stats() ServerStats

Stats 返回服务端统计快照

func (*Server) Stop

func (s *Server) Stop() error

Stop 停止服务器

type ServerConfig

type ServerConfig struct {
	// 地址
	Address string
	// 连接类型
	ConnType ConnType
	// 编解码器
	Codec Codec
	// 日志器
	Logger Logger
	// 读取缓冲区大小
	ReadBufferSize int
	// 写入缓冲区大小
	WriteBufferSize int
	// 最大连接数
	MaxConnections int
	// 读超时
	ReadTimeout time.Duration
	// 写超时
	WriteTimeout time.Duration
	// 心跳间隔(服务端:检测周期;客户端:发送周期)
	HeartbeatInterval time.Duration
	// 心跳超时(超过此时间未收到心跳则断开)
	HeartbeatTimeout time.Duration
	// 心跳 Ping 路由ID
	HeartbeatPingID uint32
	// 心跳 Pong 路由ID
	HeartbeatPongID uint32
	// 心跳 Pong 内容生成函数,入参为收到的 ping 消息,返回 pong body;为 nil 时 pong body 为空
	HeartbeatPongData func(ping Message) []byte
	// KCP 配置,ConnType 为 ConnTypeKCP 时生效
	KCPConfig *KCPConfig
	// GWS 配置,ConnType 为 ConnTypeGWS 时生效
	GWSConfig *GWSConfig

	// PoolSize 客户端连接池大小,默认 1
	PoolSize int
	// ReconnectEnable 断线后是否自动重连
	ReconnectEnable bool
	// ReconnectInitDelay 首次重连等待时间,默认 1s
	ReconnectInitDelay time.Duration
	// ReconnectMaxDelay 最大重连等待时间(指数退避上限),默认 30s
	ReconnectMaxDelay time.Duration
	// ReconnectMaxAttempts 最大重连次数,0 表示无限重试
	ReconnectMaxAttempts int
	// IDGenerator 连接 ID 生成函数,默认使用 UUID
	IDGenerator func() string
}

ServerConfig 服务器配置

func DefaultServerConfig

func DefaultServerConfig() *ServerConfig

DefaultServerConfig 返回默认服务器配置

type ServerOption

type ServerOption func(*ServerConfig)

ServerOption 服务器配置选项

func WithAddress

func WithAddress(addr string) ServerOption

WithAddress 设置地址

func WithCodec

func WithCodec(codec Codec) ServerOption

WithCodec 设置编解码器

func WithConnType

func WithConnType(t ConnType) ServerOption

WithConnType 设置连接类型

func WithGWSConfig added in v1.0.14

func WithGWSConfig(cfg *GWSConfig) ServerOption

WithGWSConfig 设置 gws 服务端配置

func WithHeartbeat

func WithHeartbeat(interval, timeout time.Duration) ServerOption

WithHeartbeat 设置心跳参数

func WithHeartbeatPongData added in v1.0.9

func WithHeartbeatPongData(fn func(ping Message) []byte) ServerOption

WithHeartbeatPongData 设置 pong 内容生成函数 fn 入参为收到的 ping 消息,返回值作为 pong 的 body

func WithHeartbeatRouteID added in v1.0.9

func WithHeartbeatRouteID(pingID, pongID uint32) ServerOption

WithHeartbeatRouteID 设置心跳包的路由ID

func WithIDGenerator added in v1.0.20

func WithIDGenerator(fn func() string) ServerOption

WithIDGenerator 设置连接 ID 生成函数

func WithKCPConfig added in v1.0.9

func WithKCPConfig(cfg *KCPConfig) ServerOption

WithKCPConfig 设置 KCP 配置

func WithLogger

func WithLogger(logger Logger) ServerOption

WithLogger 设置日志器

func WithMaxConnections

func WithMaxConnections(max int) ServerOption

WithMaxConnections 设置最大连接数

func WithPoolSize added in v1.0.19

func WithPoolSize(n int) ServerOption

WithPoolSize 设置客户端连接池大小(最小为 1)

func WithReadBufferSize

func WithReadBufferSize(size int) ServerOption

WithReadBufferSize 设置读取缓冲区大小

func WithReadTimeout

func WithReadTimeout(timeout time.Duration) ServerOption

WithReadTimeout 设置读超时

func WithReconnect added in v1.0.19

func WithReconnect(enable bool) ServerOption

WithReconnect 设置是否启用断线重连

func WithReconnectDelay added in v1.0.19

func WithReconnectDelay(initDelay, maxDelay time.Duration) ServerOption

WithReconnectDelay 设置断线重连的初始等待时间和最大等待时间

func WithReconnectMaxAttempts added in v1.0.19

func WithReconnectMaxAttempts(n int) ServerOption

WithReconnectMaxAttempts 设置最大重连次数(0 表示无限重试)

func WithWriteBufferSize

func WithWriteBufferSize(size int) ServerOption

WithWriteBufferSize 设置写入缓冲区大小

func WithWriteTimeout

func WithWriteTimeout(timeout time.Duration) ServerOption

WithWriteTimeout 设置写超时

type ServerStats added in v1.0.19

type ServerStats struct {
	// CurrentConns 当前连接数
	CurrentConns int64
	// TotalConns 累计建立的连接数
	TotalConns int64
	// TotalMessages 累计处理的消息数
	TotalMessages int64
	// TotalRecvBytes 累计接收字节数
	TotalRecvBytes int64
	// TotalSendBytes 累计发送字节数
	TotalSendBytes int64
}

ServerStats 服务端统计快照

type SimpleCodec

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

SimpleCodec 简单编解码器 包格式: [4字节totalLen(大端)] + [4字节routeID(大端)] + [body] header = 8字节,totalLen 包含自身

func NewSimpleCodec

func NewSimpleCodec(maxPacketSize ...int) *SimpleCodec

NewSimpleCodec 创建简单编解码器,maxPacketSize 默认 64KB

func (*SimpleCodec) DecodeBody added in v1.0.9

func (c *SimpleCodec) DecodeBody(routeID uint32, body []byte) (Message, error)

DecodeBody 将 body 解析为消息

func (*SimpleCodec) DecodeHeader added in v1.0.9

func (c *SimpleCodec) DecodeHeader(header []byte) (routeID uint32, bodyLen int, err error)

DecodeHeader 解析 header,返回 routeID 和 body 长度

func (*SimpleCodec) Encode

func (c *SimpleCodec) Encode(msg Message) ([]byte, error)

Encode 编码消息

func (*SimpleCodec) HeaderSize added in v1.0.9

func (c *SimpleCodec) HeaderSize() int

HeaderSize 返回固定 header 大小

func (*SimpleCodec) MaxPacketSize

func (c *SimpleCodec) MaxPacketSize() int

MaxPacketSize 返回最大包大小

type StdLogger

type StdLogger struct{}

StdLogger 使用标准库 log 的日志实现。 不再直接依赖 mlog 模块,如需使用 mlog,请在外部实现 Logger 接口并传入。

func (*StdLogger) Debug added in v1.0.9

func (l *StdLogger) Debug(msg string, args ...any)

func (*StdLogger) Error added in v1.0.9

func (l *StdLogger) Error(msg string, args ...any)

func (*StdLogger) Info added in v1.0.9

func (l *StdLogger) Info(msg string, args ...any)

func (*StdLogger) Warn added in v1.0.9

func (l *StdLogger) Warn(msg string, args ...any)

type TLVCodec

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

TLVCodec TLV格式编解码器 包格式: [1字节Type] + [2字节Length(大端)] + [Value] header = 3字节,Length 为 body 长度

func NewTLVCodec

func NewTLVCodec(maxPacketSize ...int) *TLVCodec

NewTLVCodec 创建TLV编解码器

func (*TLVCodec) DecodeBody added in v1.0.9

func (c *TLVCodec) DecodeBody(routeID uint32, body []byte) (Message, error)

DecodeBody 将 body 解析为消息

func (*TLVCodec) DecodeHeader added in v1.0.9

func (c *TLVCodec) DecodeHeader(header []byte) (routeID uint32, bodyLen int, err error)

DecodeHeader 解析 header,返回 routeID 和 body 长度

func (*TLVCodec) Encode

func (c *TLVCodec) Encode(msg Message) ([]byte, error)

Encode 编码TLV消息

func (*TLVCodec) HeaderSize added in v1.0.9

func (c *TLVCodec) HeaderSize() int

HeaderSize 返回固定 header 大小

func (*TLVCodec) MaxPacketSize

func (c *TLVCodec) MaxPacketSize() int

MaxPacketSize 返回最大包大小

type TLVMessage

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

TLVMessage TLV消息

func (*TLVMessage) Data

func (m *TLVMessage) Data() []byte

Data 返回数据

func (*TLVMessage) MsgType

func (m *TLVMessage) MsgType() byte

MsgType 返回消息类型

func (*TLVMessage) RouteID

func (m *TLVMessage) RouteID() uint32

RouteID TLVMessage 的路由ID就是 msgType

func (*TLVMessage) SetData

func (m *TLVMessage) SetData(data []byte)

SetData 设置数据

type WebSocketUpgrader

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

WebSocketUpgrader WebSocket升级器

func NewWebSocketUpgrader

func NewWebSocketUpgrader() *WebSocketUpgrader

NewWebSocketUpgrader 创建WebSocket升级器

func (*WebSocketUpgrader) SetCheckOrigin

func (u *WebSocketUpgrader) SetCheckOrigin(fn func(r *http.Request) bool)

SetCheckOrigin 设置跨域检查函数

Jump to

Keyboard shortcuts

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