// Package realtime 全站实时通信总线(一期)。 // // 设计: // - 单条 WebSocket 连接(GET /api/ws,cookie 鉴权)复用承载所有实时事件, // 后续群聊、通知、审核角标等都接入本 Hub,而不是各开连接。 // - Hub 维护连接表、房间(主题订阅)表与每用户连接计数;房间用于定向广播 // (user:{id} 个人房间、staff 管理团队房间)。 // - 所有写操作投递到 client.send 缓冲通道,由唯一的 writePump 消费, // 杜绝多 goroutine 并发写 WebSocket。 // - 在线判定 = 拥有至少 1 条活跃 WS 连接,比 last_seen_at 5 分钟启发式精确; // 上线/下线事件只推送给 staff 房间。 package realtime import ( "encoding/json/v2" "sync" ) // 事件类型:前后端共享的单层 JSON 协议契约,禁止手拼 JSON / 双重编码 const ( EventHello = "hello" // 建连欢迎帧(含在线用户快照) EventPong = "pong" // 应用层心跳应答 EventSettingsChanged = "settings:changed" // 站点设置(主题色)变更,全员广播 EventPresenceUpdate = "presence:update" // 用户上/下线,仅 staff 房间 EventChatMessage = "chat:message" // 群聊新消息,推 chat:{roomID} EventChatMembership = "chat:membership" // 群成员关系变更(被踢/群解散/被邀请),推 user:{id} EventNotificationNew = "notification:new" // 新站内通知(如群聊 @),推 user:{id} EventFeedChanged = "feed:changed" // 帖子流有公开新内容(三期),全员广播 EventChatRecalled = "chat:message_recalled" // 消息撤回,推 chat:{roomID} EventChatUnread = "chat:unread" // 未读增量,推 user:{id}(未订阅房间也能实时角标) EventModerationChanged = "moderation:changed" // 待审队列变化,推 staff 房间(客户端各自 HTTP 校准角标) EventSessionReplaced = "session:replaced" // 登录态被顶/被剔除,推 user:{id}(旧设备即时感知并提示重新登录) ) // RoomStaff 管理团队房间(板块管理员及以上) const RoomStaff = "staff" // RoomChat 群聊房间(仅群成员可通过 join 帧订阅) func RoomChat(roomID uint) string { return "chat:" + uintToString(roomID) } func roomUser(userID uint) string { return "user:" + uintToString(userID) } // RoomUser 个人房间(导出供按用户定向推送) func RoomUser(userID uint) string { return roomUser(userID) } // Envelope 实时消息统一信封 type Envelope struct { Type string `json:"type"` Data interface{} `json:"data,omitempty"` } // Hub 连接与房间注册中心(单实例内存版;多实例时由 PG LISTEN/NOTIFY 桥接) type Hub struct { mu sync.RWMutex // 全部活跃连接 clients map[*Client]struct{} // 房间 -> 连接集合 rooms map[string]map[*Client]struct{} // 用户 -> 活跃连接数(同一用户多标签页) userConnCnt map[uint]int } func NewHub() *Hub { return &Hub{ clients: make(map[*Client]struct{}), rooms: make(map[string]map[*Client]struct{}), userConnCnt: make(map[uint]int), } } // register 注册新连接。返回 cameOnline=true 表示该用户 0→1 首次上线。 func (h *Hub) register(c *Client) (cameOnline bool) { h.mu.Lock() h.clients[c] = struct{}{} n := h.userConnCnt[c.userID] h.userConnCnt[c.userID] = n + 1 h.joinLocked(c, roomUser(c.userID)) if c.staff { h.joinLocked(c, RoomStaff) } h.mu.Unlock() return n == 0 } // unregister 注销连接,返回 wentOffline=true 表示该用户最后一条连接断开。 func (h *Hub) unregister(c *Client) (wentOffline bool) { h.mu.Lock() if _, ok := h.clients[c]; !ok { h.mu.Unlock() return false } delete(h.clients, c) for room := range c.rooms { h.leaveLocked(c, room) } clear(c.rooms) n := h.userConnCnt[c.userID] - 1 if n <= 0 { delete(h.userConnCnt, c.userID) wentOffline = true } else { h.userConnCnt[c.userID] = n } h.mu.Unlock() return wentOffline } func (h *Hub) joinLocked(c *Client, room string) { m := h.rooms[room] if m == nil { m = make(map[*Client]struct{}) h.rooms[room] = m } m[c] = struct{}{} c.rooms[room] = struct{}{} } func (h *Hub) leaveLocked(c *Client, room string) { if m := h.rooms[room]; m != nil { delete(m, c) if len(m) == 0 { delete(h.rooms, room) } } delete(c.rooms, room) } // BroadcastAll 向全部连接广播 func (h *Hub) BroadcastAll(env Envelope) { h.mu.RLock() clients := make([]*Client, 0, len(h.clients)) for c := range h.clients { clients = append(clients, c) } h.mu.RUnlock() h.dispatch(clients, env) } // BroadcastRoom 向指定房间广播 func (h *Hub) BroadcastRoom(room string, env Envelope) { h.mu.RLock() m := h.rooms[room] clients := make([]*Client, 0, len(m)) for c := range m { clients = append(clients, c) } h.mu.RUnlock() h.dispatch(clients, env) } // JoinClientRoom 连接动态加入房间(群聊订阅;调用方必须先完成成员鉴权)。 // 返回 false 表示连接已注销。 func (h *Hub) JoinClientRoom(c *Client, room string) bool { h.mu.Lock() defer h.mu.Unlock() if _, ok := h.clients[c]; !ok { return false } h.joinLocked(c, room) return true } // LeaveClientRoom 连接动态离开房间 func (h *Hub) LeaveClientRoom(c *Client, room string) { h.mu.Lock() defer h.mu.Unlock() h.leaveLocked(c, room) } // RemoveUserRoom 把某用户的全部活跃连接移出房间(群主踢人后立即停止实时投递) func (h *Hub) RemoveUserRoom(userID uint, room string) { h.mu.Lock() var targets []*Client for c := range h.rooms[room] { if c.userID == userID { targets = append(targets, c) } } for _, c := range targets { h.leaveLocked(c, room) } h.mu.Unlock() } // BroadcastUser 向某用户的全部活跃连接(多标签页)单播 func (h *Hub) BroadcastUser(userID uint, env Envelope) { h.BroadcastRoom(roomUser(userID), env) } // BroadcastPresence 上/下线事件推送给管理团队 func (h *Hub) BroadcastPresence(userID uint, online bool) { h.BroadcastRoom(RoomStaff, Envelope{ Type: EventPresenceUpdate, Data: map[string]interface{}{ "user_id": userID, "online": online, }, }) } // OnlineUserIDs 当前在线用户快照(供 staff 建连时校正绿点) func (h *Hub) OnlineUserIDs() []uint { h.mu.RLock() defer h.mu.RUnlock() ids := make([]uint, 0, len(h.userConnCnt)) for id := range h.userConnCnt { ids = append(ids, id) } return ids } func (h *Hub) dispatch(clients []*Client, env Envelope) { if len(clients) == 0 { return } payload, err := json.Marshal(env) if err != nil { return } for _, c := range clients { select { case c.send <- payload: default: // 慢消费者:缓冲已满直接关闭连接,等客户端指数退避重连补齐 c.forceClose() } } }