netutil

package
v1.1.1 Latest Latest
Warning

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

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

Documentation

Index

Constants

View Source
const (
	HeartbeatInterval = 60 * time.Second // 心跳发送间隔
	TaskQueueBufSize  = 500              // 连接任务队列缓冲大小(ants池用)
	ConnPoolWorkers   = 1                // 每个连接任务池的协程数-只能为1
	DefaultMaxRetries = 3                // 默认最大重试次数

	DefaultRetryCnts int = 3
)

-------------------------- 配置常量 --------------------------

View Source
const (
	TextType   int = websocket.TextMessage
	BinaryType int = websocket.BinaryMessage
)

消息类型

Variables

This section is empty.

Functions

This section is empty.

Types

type ClientInfo

type ClientInfo struct {
	ConnectionID string    // 业务用户ID
	ConnTime     time.Time //连接存活时长

}

ClientInfo 客户端信息

func (*ClientInfo) GetID

func (c *ClientInfo) GetID() string

type CloseTask

type CloseTask struct {
	Connection *WsConnection
	Reason     string
	Done       chan struct{}
}

连接关闭任务

func (*CloseTask) Execute

func (t *CloseTask) Execute()

type HandleMsgTask

type HandleMsgTask struct {
	Connection  *WsConnection //回指向所属连接
	MessageType int           //消息类型
	Body        []byte        //消息载荷
	Retries     int           // 当前重试次数
}

HandleMsgTask 消息处理任务(实现Task接口)

func (*HandleMsgTask) Execute

func (t *HandleMsgTask) Execute()

type HttpOption

type HttpOption = func(h *HttpServer)

func WithCorsConfig

func WithCorsConfig(mode string, cfg *conf.LogConfig) HttpOption

配置 CORS 中间件并重定向日志输出器

func WithHttpPort

func WithHttpPort(port int) HttpOption

func WithInitHttpLog

func WithInitHttpLog(outlogger NetLogger) HttpOption

type HttpServer

type HttpServer struct {
	Server *gin.Engine
	Port   int
}

func NewHttpServer

func NewHttpServer(opts ...HttpOption) *HttpServer

func (*HttpServer) GET added in v1.0.9

func (s *HttpServer) GET(url string, handler gin.HandlerFunc)

func (*HttpServer) GinEngine

func (s *HttpServer) GinEngine() *gin.Engine

获取gin引擎

func (*HttpServer) Put

func (s *HttpServer) Put(url string, handler gin.HandlerFunc)

func (*HttpServer) Start

func (s *HttpServer) Start() error

type ManagerOption

type ManagerOption = func(m *WsConnectionManager)

func WithCloseHandler

func WithCloseHandler(handler OnCloseCallBack) ManagerOption

func WithMessageHandler

func WithMessageHandler(handler OnMessageCallBack) ManagerOption

type NetLogger

type NetLogger interface {
	Debug(expand *map[string]any, format string)
	Info(expand *map[string]any, format string)
	Error(expand *map[string]any, format string)
}

日志器接口

var Log NetLogger

type OnCloseCallBack

type OnCloseCallBack func(*WsConnection, string)

连接关闭时触发

type OnMessageCallBack

type OnMessageCallBack func(*WsConnection, int, []byte) error

消息处理函数类型 返回error供业务层判断处理结

type PingTask

type PingTask struct {
	Connection *WsConnection
	Retries    int // 当前重试次数
}

HeartbeatTask 心跳任务(实现Task接口)

func (*PingTask) Execute

func (t *PingTask) Execute()

type SendMsgTask

type SendMsgTask struct {
	Connection  *WsConnection //回指向所属连接
	MessageType int           //消息类型
	Body        []byte        //消息载荷
	Done        chan error    // 同步返回发送结果
	Retries     int           // 当前重试次数
}

SendMsgTask 消息发送任务(实现Task接口)

func (*SendMsgTask) Execute

func (t *SendMsgTask) Execute()

任务执行接口

type Task

type Task interface {
	Execute() // 执行任务
}

-------------------------- 任务抽象 -------------------------- Task 任务接口:所有连接操作需实现此接口

type WebSocketServer

type WebSocketServer struct {
	Server        *HttpServer
	WsConnManager *WsConnectionManager
	WsUp          websocket.Upgrader
}

func NewDefaultWsServer

func NewDefaultWsServer(http *HttpServer) *WebSocketServer

创建默认的websocker-server

func (*WebSocketServer) SetExceptionHandler

func (w *WebSocketServer) SetExceptionHandler(on_close OnCloseCallBack)

func (*WebSocketServer) SetMessageHandler

func (w *WebSocketServer) SetMessageHandler(on_message OnMessageCallBack)

设置事件发生时触发的回调方法

func (*WebSocketServer) Start

func (w *WebSocketServer) Start() error

启动websocker服务器

type WsConnection

type WsConnection struct {
	Connection     *websocket.Conn      // 底层WebSocket连接
	PeerInfo       ClientInfo           // 客户端信息
	Owner          *WsConnectionManager // 关联的连接管理器
	MessageFunc    OnMessageCallBack    // 消息处理回调函数
	CloseFunc      OnCloseCallBack      // 关闭回调函数
	HeartbeatTimer *time.Timer          // 心跳定时器
	TaskPool       *ants.Pool           // 任务池
	IsClosed       atomic.Bool          // 连接关闭状态
	Wg             sync.WaitGroup       // 等待任务完成
}

连接ID == clientinfo的ID

func (*WsConnection) Close

func (conn *WsConnection) Close(reason string)

对外关闭连接接口--资源完整清理

func (*WsConnection) GetID

func (conn *WsConnection) GetID() string

GetID 获取连接唯一标识

func (*WsConnection) Recv

func (conn *WsConnection) Recv() (messageType int, p []byte, err error)

读消息

func (*WsConnection) Send

func (conn *WsConnection) Send(messageType int, body []byte) error

发送消息

func (*WsConnection) SendMsg

func (conn *WsConnection) SendMsg(ctx context.Context, msgType int, msg []byte) error

对外发送消息接口(协程安全+超时控制)

func (*WsConnection) SubmitTask

func (conn *WsConnection) SubmitTask(task Task) error

-------------------------- WsConnection方法 -------------------------- 任务提交

type WsConnectionManager

type WsConnectionManager struct {
	IDMapConnection map[string]*WsConnection
	ConnectionMapID map[*WsConnection]ClientInfo
	Locker          sync.RWMutex
	MessageFunc     OnMessageCallBack // 消息处理回调函数
	CloseFunc       OnCloseCallBack   // 关闭回调函数
}

连接管理器

func GetGlobalMgr

func GetGlobalMgr(opts ...ManagerOption) *WsConnectionManager

创建连接管理器

func (*WsConnectionManager) AddConnectionToManager

func (m *WsConnectionManager) AddConnectionToManager(ws *websocket.Conn, peer ClientInfo) error

注册新连接到管理器mgr

func (*WsConnectionManager) FindConnByPeerID

func (m *WsConnectionManager) FindConnByPeerID(conn_id string) (*WsConnection, error)

-------------------------- 连接查找与移除 -------------------------- 根据peer查找连接conn

func (*WsConnectionManager) NewConnection

func (m *WsConnectionManager) NewConnection(ws *websocket.Conn, peer ClientInfo) (*WsConnection, error)

-------------------------- 连接创建与注册 -------------------------- 创建单个Connection实例-内部使用

func (*WsConnectionManager) RemoveConn

func (m *WsConnectionManager) RemoveConn(conn *WsConnection) error

RemoveConn 从管理器移除连接

func (*WsConnectionManager) SetOnCloseCb

func (m *WsConnectionManager) SetOnCloseCb(cb OnCloseCallBack)

func (*WsConnectionManager) SetOnMessageCb

func (m *WsConnectionManager) SetOnMessageCb(cb OnMessageCallBack)

type WsOption

type WsOption = func(s *WebSocketServer)

func WithInitDefaultWsUpgrader

func WithInitDefaultWsUpgrader(checkorigin bool) WsOption

初始化默认的ws协议升级器

func WithInitWsServerLog

func WithInitWsServerLog(outlogger NetLogger) WsOption

注入日志器

Directories

Path Synopsis
client command

Jump to

Keyboard shortcuts

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