diff --git a/controller/websocket_controller.go b/controller/websocket_controller.go index 94a4820..8053626 100644 --- a/controller/websocket_controller.go +++ b/controller/websocket_controller.go @@ -1,3 +1,8 @@ +/** + * package controller + * 作用:处理核心业务逻辑,包括 WebSocket 连接管理、消息转发、Redis 集群通信以及 API 接口实现。 + * 该包是连接网络层(Gin)与数据层(Gorm/Redis)的桥梁。 + */ package controller import ( @@ -27,7 +32,15 @@ var upgrader = websocket.Upgrader{ }, } -// WebSocketController WebSocket控制器结构体 +/** + * WebSocketController + * 结构体定义:WebSocket 服务的主控制器。 + * 包含: + * - 客户端连接池 (Clients) + * - 读写锁 (ClientsMux) 保护连接池安全 + * - Redis 客户端 (RedisCli) 用于集群通信 + * - 数据库连接 (DB) + */ type WebSocketController struct { Clients map[string]*websocket.Conn ClientsMux sync.RWMutex @@ -40,6 +53,11 @@ type WebSocketController struct { DB *gorm.DB // 数据库连接 } +/** + * NewWebSocketController + * 功能:创建并初始化 WebSocketController 实例。 + * @param db *gorm.DB 数据库连接实例 + */ func NewWebSocketController(db *gorm.DB) *WebSocketController { return &WebSocketController{ Clients: make(map[string]*websocket.Conn), @@ -48,7 +66,11 @@ func NewWebSocketController(db *gorm.DB) *WebSocketController { } } -// ConfigureSystem 配置系统参数 +/** + * ConfigureSystem + * 功能:加载系统配置,如 NodeID 和端口,并设置日志格式。 + * 也会启动节点注册协程。 + */ func (c *WebSocketController) ConfigureSystem() { c.NodeID = utils.GetEnv("NODE_ID", "local") @@ -69,7 +91,11 @@ func (c *WebSocketController) ConfigureSystem() { go c.registerNode() // 现在 RedisCli 已初始化 } -// 注册节点到Redis +/** + * registerNode + * 功能:定时将当前节点的 IP 和端口注册到 Redis 中,以便其他服务发现。 + * 这是一个心跳机制,每 20 秒执行一次。 + */ func (c *WebSocketController) registerNode() { // 添加空值检查 if c.RedisCli == nil { @@ -91,7 +117,10 @@ func (c *WebSocketController) registerNode() { } } -// 配置日志记录器 +/** + * configureLogger + * 功能:初始化日志目录和写入器配置。 + */ func (c *WebSocketController) configureLogger() { logDir := utils.GetEnv("LOG_DIR", "./logs") if err := os.MkdirAll(logDir, 0755); err != nil { @@ -107,12 +136,18 @@ func (c *WebSocketController) configureLogger() { log.SetOutput(c.getLogWriter(logDir, logFreq)) } -// 获取日志文件写入器 +/** + * getLogWriter + * 功能:工厂方法,获取 DailyFileWriter 实例。 + */ func (c *WebSocketController) getLogWriter(logDir, freq string) *utils.DailyFileWriter { return utils.NewDailyFileWriter(logDir, freq, c.Logger, c.NodeID) } -// InitRedisClient 初始化Redis客户端 +/** + * InitRedisClient + * 功能:建立 Redis 连接并验证连通性。 + */ func (c *WebSocketController) InitRedisClient() { redisAddr := utils.GetEnv("REDIS_ADDR", "localhost:6379") redisPassword := utils.GetEnv("REDIS_PASSWORD", "") @@ -133,7 +168,10 @@ func (c *WebSocketController) checkRedisConnection() error { return err } -// PrintStartupInfo 打印启动信息 +/** + * PrintStartupInfo + * 功能:在控制台打印详细的服务启动信息,包括节点ID、端口、Redis状态及支持的消息类型。 + */ func (c *WebSocketController) PrintStartupInfo() { hostname, _ := os.Hostname() log.Printf("🚀🚀🚀🚀🚀🚀🚀🚀 WebSocket服务启动: NodeID=%s", c.NodeID) @@ -148,7 +186,11 @@ func (c *WebSocketController) PrintStartupInfo() { log.Println("🔗🔗🔗🔗🔗🔗🔗🔗 等待客户端连接...") } -// HealthHandler 健康检查处理器 +/** + * HealthHandler + * 功能:API 接口,用于健康检查。 + * 路径:GET /api/health + */ func (c *WebSocketController) HealthHandler(ctx *gin.Context) { ctx.JSON(http.StatusOK, gin.H{ "status": "ok", @@ -158,7 +200,11 @@ func (c *WebSocketController) HealthHandler(ctx *gin.Context) { }) } -// IsUserOnline 检查某个用户ID是否在线 +/** + * IsUserOnline + * 功能:API 接口,检查指定用户当前是否在线(通过 Redis 查询)。 + * 路径:GET /api/check-user-online + */ func (c *WebSocketController) IsUserOnline(ctx *gin.Context) { // 优先尝试查询参数,再尝试路径参数 userID := ctx.Query("user_id") @@ -200,7 +246,11 @@ func (c *WebSocketController) IsUserOnline(ctx *gin.Context) { }) } -// HandleWebSocket WebSocket连接处理器 +/** + * HandleWebSocket + * 功能:处理 WebSocket 握手升级,并开启消息监听循环。 + * 路径:GET /ws + */ func (c *WebSocketController) HandleWebSocket(ctx *gin.Context) { start := time.Now() clientIP := ctx.ClientIP() @@ -248,7 +298,12 @@ func (c *WebSocketController) HandleWebSocket(ctx *gin.Context) { } } -// 处理客户端消息 +/** + * handleClientMessage + * 功能:处理从客户端接收到的消息,支持绑定、发消息、通话信令等操作。 + * @param senderID string 发送者的客户端ID + * @param message []byte 原始消息体 + */ func (c *WebSocketController) handleClientMessage(senderID string, message []byte) { log.Printf("📥📥📥📥📥📥📥📥 收到客户端消息: \n SenderID=%s \n Size=%d bytes", senderID, len(message)) @@ -361,7 +416,10 @@ func (c *WebSocketController) handleClientMessage(senderID string, message []byt log.Printf("⚠️ 无法识别的消息格式: \n Size=%d bytes \n Message=%s", len(message), string(message)) } -// 处理通话邀请 +/** + * handleCallInvite + * 功能:处理WebRTC通话邀请信令。 + */ func (c *WebSocketController) handleCallInvite(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📞 处理通话邀请: \n CallID=%s \n CallType=%d \n From=%s \n To=%s", signal.CallID, signal.CallType, senderUserID, signal.CalleeID) @@ -385,7 +443,10 @@ func (c *WebSocketController) handleCallInvite(senderID, senderUserID string, si } } -// 处理通话接受 +/** + * handleCallAccept + * 功能:处理被叫方接受通话的信令。 + */ func (c *WebSocketController) handleCallAccept(senderID, senderUserID string, signal models.CallSignal) { log.Printf("✅ 处理通话接受: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CallerID) @@ -408,7 +469,10 @@ func (c *WebSocketController) handleCallAccept(senderID, senderUserID string, si } } -// 处理通话拒绝 +/** + * handleCallReject + * 功能:处理被叫方拒绝通话的信令。 + */ func (c *WebSocketController) handleCallReject(senderID, senderUserID string, signal models.CallSignal) { log.Printf("❌ 处理通话拒绝: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CallerID) @@ -431,7 +495,10 @@ func (c *WebSocketController) handleCallReject(senderID, senderUserID string, si } } -// 处理通话结束 +/** + * handleCallEnd + * 功能:处理通话结束/挂断信令。 + */ func (c *WebSocketController) handleCallEnd(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📞 处理通话结束: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID) @@ -454,7 +521,10 @@ func (c *WebSocketController) handleCallEnd(senderID, senderUserID string, signa } } -// 处理Offer信令 +/** + * handleCallOffer + * 功能:处理 WebRTC 的 Offer 信令(SDP交换)。 + */ func (c *WebSocketController) handleCallOffer(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📞 处理Offer信令: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID) @@ -477,7 +547,10 @@ func (c *WebSocketController) handleCallOffer(senderID, senderUserID string, sig } } -// 处理Answer信令 +/** + * handleCallAnswer + * 功能:处理 WebRTC 的 Answer 信令(SDP交换)。 + */ func (c *WebSocketController) handleCallAnswer(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📞 处理Answer信令: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CallerID) @@ -500,7 +573,10 @@ func (c *WebSocketController) handleCallAnswer(senderID, senderUserID string, si } } -// 处理通话候选 +/** + * handleCallCandidate + * 功能:处理 ICE Candidate 信令(网络协商)。 + */ func (c *WebSocketController) handleCallCandidate(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📶 处理通话候选: \n CallID=%s From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID) @@ -523,7 +599,10 @@ func (c *WebSocketController) handleCallCandidate(senderID, senderUserID string, } } -// 处理通话挂断 +/** + * handleCallHangup + * 功能:处理主动挂断通话。 + */ func (c *WebSocketController) handleCallHangup(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📞 处理通话挂断: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID) @@ -546,7 +625,10 @@ func (c *WebSocketController) handleCallHangup(senderID, senderUserID string, si } } -// 处理通话掉线 +/** + * handleCallDisconnected + * 功能:处理连接意外断开。 + */ func (c *WebSocketController) handleCallDisconnected(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📞 处理通话掉线: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID) @@ -569,7 +651,10 @@ func (c *WebSocketController) handleCallDisconnected(senderID, senderUserID stri } } -// 处理通话终止 +/** + * handleCallTerminated + * 功能:处理通话终止信号。 + */ func (c *WebSocketController) handleCallTerminated(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📞 处理通话终止: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID) @@ -592,7 +677,10 @@ func (c *WebSocketController) handleCallTerminated(senderID, senderUserID string } } -// 处理无人接听 +/** + * handleCallNoAnswer + * 功能:处理无人接听状态。 + */ func (c *WebSocketController) handleCallNoAnswer(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📞 处理无人接听: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CallerID) @@ -615,7 +703,10 @@ func (c *WebSocketController) handleCallNoAnswer(senderID, senderUserID string, } } -// 处理忙线状态 +/** + * handleCallBusy + * 功能:处理忙线状态。 + */ func (c *WebSocketController) handleCallBusy(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📞 处理忙线状态: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CallerID) @@ -638,7 +729,10 @@ func (c *WebSocketController) handleCallBusy(senderID, senderUserID string, sign } } -// 处理通话失败 +/** + * handleCallFailed + * 功能:处理通话建立失败状态。 + */ func (c *WebSocketController) handleCallFailed(senderID, senderUserID string, signal models.CallSignal) { log.Printf("📞 处理通话失败: \n CallID=%s \n From=%s \n To=%s", signal.CallID, senderUserID, signal.CalleeID) @@ -661,7 +755,11 @@ func (c *WebSocketController) handleCallFailed(senderID, senderUserID string, si } } -// SendMessageHandler API消息发送处理器 +/** + * SendMessageHandler + * 功能:API 接口,用于服务器端主动发送消息(HTTP -> WebSocket)。 + * 路径:POST /api/send + */ func (c *WebSocketController) SendMessageHandler(ctx *gin.Context) { start := time.Now() @@ -724,7 +822,11 @@ func (c *WebSocketController) SendMessageHandler(ctx *gin.Context) { ctx.JSON(http.StatusBadRequest, gin.H{"error": "不支持的request_type"}) } -// 用户绑定处理器 +/** + * BindHandler + * 功能:API 接口,用于强制绑定 ClientID 和 UserID。 + * 路径:POST /api/bind + */ func (c *WebSocketController) BindHandler(ctx *gin.Context) { start := time.Now() @@ -758,7 +860,11 @@ func (c *WebSocketController) BindHandler(ctx *gin.Context) { req.UserID, req.ClientID, time.Since(start)) } -// SendToUserHandler 通过用户ID发送消息处理器(支持通话信令) +/** + * SendToUserHandler + * 功能:API 接口,向指定用户的所有在线设备发送消息,并持久化到数据库。 + * 路径:POST /api/send-to-user + */ // SendToUserHandler 通过用户ID发送消息处理器(支持通话信令) func (c *WebSocketController) SendToUserHandler(ctx *gin.Context) { start := time.Now() @@ -841,7 +947,11 @@ func (c *WebSocketController) SendToUserHandler(ctx *gin.Context) { log.Printf("✅ API用户发送请求完成: Duration=%s", time.Since(start)) } -// 获取聊天记录(分页) +/** + * GetMessagesHandler + * 功能:API 接口,分页查询历史聊天记录。 + * 路径:GET /api/messages + */ func (c *WebSocketController) GetMessagesHandler(ctx *gin.Context) { // 解析查询参数 roomId := ctx.Query("room_id") @@ -892,7 +1002,11 @@ func (c *WebSocketController) GetMessagesHandler(ctx *gin.Context) { }) } -// 同步聊天记录(全部) +/** + * SyncMessagesHandler + * 功能:API 接口,获取全部聊天记录(用于同步)。 + * 路径:GET /api/messages/sync + */ func (c *WebSocketController) SyncMessagesHandler(ctx *gin.Context) { roomId := ctx.Query("room_id") senderUserId := ctx.Query("sender_user_id") @@ -921,7 +1035,11 @@ func (c *WebSocketController) SyncMessagesHandler(ctx *gin.Context) { }) } -// 添加客户端连接 +/** + * addClient + * 功能:将新的 WebSocket 连接添加到本地连接池。 + * 注意:使用 Mutex 锁保证线程安全。 + */ func (c *WebSocketController) addClient(clientID string, conn *websocket.Conn) { c.ClientsMux.Lock() defer c.ClientsMux.Unlock() @@ -929,7 +1047,10 @@ func (c *WebSocketController) addClient(clientID string, conn *websocket.Conn) { log.Printf("📊📊📊📊📊📊📊📊 当前连接数: %d", len(c.Clients)) } -// 移除客户端连接 +/** + * removeClient + * 功能:从本地连接池移除客户端,并清理 Redis 中的绑定关系。 + */ func (c *WebSocketController) removeClient(clientID string) { c.ClientsMux.Lock() defer c.ClientsMux.Unlock() @@ -943,7 +1064,10 @@ func (c *WebSocketController) removeClient(clientID string) { } } -// 清理用户绑定关系 +/** + * cleanupUserBinding + * 功能:从 Redis 中清除用户与客户端ID的映射关系。 + */ func (c *WebSocketController) cleanupUserBinding(clientID string) { // 获取用户ID userID, err := c.getUserIDByClientID(clientID) @@ -962,7 +1086,12 @@ func (c *WebSocketController) cleanupUserBinding(clientID string) { log.Printf("🧹🧹🧹🧹🧹🧹🧹🧹 清理绑定关系: ClientID=%s | UserID=%s", clientID, userID) } -// 绑定客户端到用户 +/** + * bindClientToUser + * 功能:在 Redis 中建立 UserID 和 ClientID 的双向映射。 + * 映射1:client_user_mapping:clientID -> userID (string) + * 映射2:user_client_mapping:userID -> [clientID1, clientID2] (set) + */ func (c *WebSocketController) bindClientToUser(clientID, userID string) error { // 设置客户端->用户映射 clientUserKey := fmt.Sprintf("%s:%s", models.ClientUserKey, clientID) @@ -1000,7 +1129,11 @@ func (c *WebSocketController) getClientIDsByUserID(userID string) ([]string, err return clientIDs, nil } -// 发送消息给指定客户端 +/** + * sendMessageToClient + * 功能:向指定 ClientID 发送消息。 + * 逻辑:如果 Client 在本节点,直接通过 WebSocket 发送;如果不在,尝试通过 Redis 广播转发。 + */ func (c *WebSocketController) sendMessageToClient(clientID string, message models.ClientReceivedMessage) error { c.ClientsMux.RLock() conn, exists := c.Clients[clientID] @@ -1025,7 +1158,10 @@ func (c *WebSocketController) sendMessageToClient(clientID string, message model return nil } -// 发送消息给指定用户 +/** + * sendMessageToUser + * 功能:向指定用户的所有设备发送消息。 + */ func (c *WebSocketController) sendMessageToUser(userID string, message models.ClientReceivedMessage) error { clientIDs, err := c.getClientIDsByUserID(userID) if err != nil { @@ -1061,7 +1197,10 @@ func (c *WebSocketController) sendMessageToUser(userID string, message models.Cl return nil } -// 转发消息到其他节点 +/** + * forwardMessageToOtherNodes + * 功能:当目标客户端不在本机时,通过 Redis Pub/Sub 将消息广播到集群其他节点。 + */ func (c *WebSocketController) forwardMessageToOtherNodes(clientID string, message models.ClientReceivedMessage) error { redisMsg := models.RedisMessage{ SenderNodeID: c.NodeID, @@ -1089,7 +1228,10 @@ func (c *WebSocketController) forwardMessageToOtherNodes(clientID string, messag return nil } -// 启动Redis消息监听 +/** + * SubscribeToRedis + * 功能:订阅 Redis 广播频道,监听来自其他节点的消息转发请求。 + */ func (c *WebSocketController) SubscribeToRedis() { go func() { pattern := "websocket:forward:*" diff --git a/main.go b/main.go index 23da2d4..96a77f5 100644 --- a/main.go +++ b/main.go @@ -15,6 +15,11 @@ import ( "github.com/gin-gonic/gin" ) +/** + * ensureLogDir + * 功能:检查指定的日志目录是否存在,如果不存在则尝试创建该目录。 + * @param logDir string 日志目录路径 + */ func ensureLogDir(logDir string) { if _, err := os.Stat(logDir); os.IsNotExist(err) { log.Printf("📂 日志目录不存在,正在创建: %s", logDir) @@ -25,6 +30,14 @@ func ensureLogDir(logDir string) { } } +/** + * setupLogger + * 功能:配置全局日志系统。 + * 逻辑: + * 1. 读取环境变量配置日志目录和轮转频率。 + * 2. 创建自定义的 DailyFileWriter 以支持按天/按小时分割日志。 + * 3. 将标准 log 的输出重定向到文件系统。 + */ func setupLogger() { // 设置日志目录(优先级:环境变量 > 默认值) logDir := os.Getenv("LOG_DIR") @@ -49,7 +62,11 @@ func setupLogger() { log.SetOutput(fileWriter) } -// 初始化数据库连接 +/** + * initDB + * 功能:初始化 MySQL 数据库连接。 + * 返回:*gorm.DB 数据库连接实例 + */ func initDB() *gorm.DB { dsn := utils.GetEnv("DB_DSN", "root:password@tcp(localhost:3306)/xk_chat?charset=utf8mb4&parseTime=True&loc=Local") db, err := gorm.Open(mysql.Open(dsn), &gorm.Config{}) @@ -66,6 +83,17 @@ func initDB() *gorm.DB { return db } +/** + * main + * 功能:程序主入口。 + * 流程: + * 1. 初始化日志。 + * 2. 解析命令行参数(端口、节点ID)。 + * 3. 初始化数据库和WebSocket控制器。 + * 4. 启动 Redis 订阅(用于集群消息同步)。 + * 5. 注册 HTTP/WebSocket 路由。 + * 6. 启动 Web Server。 + */ func main() { // 初始化日志系统(必须放在最先) setupLogger() diff --git a/models/messages.go b/models/messages.go index 196679d..8bce3b1 100644 --- a/models/messages.go +++ b/models/messages.go @@ -1,7 +1,16 @@ +/** + * package models + * 作用:定义项目中所有的数据结构,包括 HTTP 请求/响应载荷、数据库模型、Redis 消息格式以及常量定义。 + */ package models import "time" +/** + * SendMessagePayload + * 结构体:客户端发送消息的请求载荷。 + * 用途:用于解析 WebSocket 或 HTTP POST 请求体中的 JSON 数据。 + */ // 统一消息结构(发送方使用) type SendMessagePayload struct { RequestType string `json:"request_type"` // 消息类型 @@ -17,6 +26,10 @@ type SendMessagePayload struct { CallStatus string `json:"call_status,omitempty"` // 通话状态:invite/accepted/rejected/ended/candidate } +/** + * SendToUserPayload + * 结构体:仅通过用户ID发送消息的简化载荷。 + */ // 通过用户ID发送消息的结构 type SendToUserPayload struct { RoomId string `json:"room_id"` // 发送者用户ID @@ -30,6 +43,11 @@ type SendToUserPayload struct { CallStatus string `json:"call_status,omitempty"` // 通话状态:invite/accepted/rejected/ended/candidate/offer/answer } +/** + * XkChatMessage + * 结构体:数据库模型,对应 xk_chat_messages 表(假设)。 + * 用途:持久化存储聊天记录。 + */ // 聊天消息记录模型 type XkChatMessage struct { ID uint `gorm:"primaryKey" json:"id"` @@ -44,6 +62,11 @@ type XkChatMessage struct { CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"` // 创建时间 } +/** + * ClientReceivedMessage + * 结构体:客户端最终接收到的消息格式。 + * 用途:服务端推送到客户端的 JSON 结构。 + */ // 客户端接收消息结构(最终格式) type ClientReceivedMessage struct { ID int `json:"id"` // 房间号 @@ -62,6 +85,11 @@ type ClientReceivedMessage struct { CallStatus string `json:"call_status,omitempty"` // 通话状态:invite/accepted/rejected/ended/candidate/offer/answer } +/** + * RedisMessage + * 结构体:Redis Pub/Sub 消息结构。 + * 用途:用于跨节点消息转发。 + */ // Redis 消息结构 type RedisMessage struct { SenderNodeID string `json:"sender_node_id"` // 发送节点ID @@ -75,6 +103,10 @@ type BindRequest struct { ClientID string `json:"client_id"` // 客户端ID } +/** + * CallSignal + * 结构体:WebRTC 信令数据载荷。 + */ // 通话信令结构 type CallSignal struct { CallStatus string `json:"call_status"` // 动作: invite/accept/reject/end/candidate/offer/answer diff --git a/route/route.go b/route/route.go index d9602ad..95db3a3 100644 --- a/route/route.go +++ b/route/route.go @@ -1,3 +1,7 @@ +/** + * package route + * 作用:定义 Gin Web Server 的路由规则。 + */ package route import ( @@ -6,6 +10,12 @@ import ( "github.com/gin-gonic/gin" ) +/** + * SetupRoutes + * 功能:配置 HTTP 和 WebSocket 路由。 + * @param router *gin.Engine Gin引擎实例 + * @param wsCtrl *controller.WebSocketController 控制器实例 + */ func SetupRoutes(router *gin.Engine, wsCtrl *controller.WebSocketController) { router.Use(gin.Recovery())