feat: 实时消息与管理端运营能力

新增群聊/私信与公共 WebSocket 总线,站长用户内容档案,并将通知与待审角标改为推送驱动;同步精简管理仪表盘并修复主题换肤 DOM 冲突。

Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
2026-09-15 22:58:16 +08:00
parent 29db9ae9b5
commit 42c395a864
55 changed files with 7607 additions and 346 deletions

206
backend/realtime/client.go Normal file
View File

@@ -0,0 +1,206 @@
package realtime
import (
"encoding/json/v2"
"strconv"
"strings"
"sync/atomic"
"time"
"github.com/gorilla/websocket"
)
func uintToString(n uint) string {
return strconv.FormatUint(uint64(n), 10)
}
const (
// writeWait 单次写帧超时
writeWait = 10 * time.Second
// pongWait 无 Pong 超时:浏览器自动响应协议层 Ping,超时即判定断连
pongWait = 70 * time.Second
// pingPeriod 协议层 Ping 周期(须小于 pongWait)
pingPeriod = 30 * time.Second
// maxMessageSize 仅接收应用层心跳,限制 4KiB 足够
maxMessageSize = 4 << 10
// sendBufferSize 慢消费保护阈值
sendBufferSize = 32
// CloseAuthInvalid 鉴权失效(封禁/强制下线)自定义关闭码,前端不自动重连
CloseAuthInvalid = 4401
)
// Client 一条 WebSocket 连接
type Client struct {
hub *Hub
conn *websocket.Conn
userID uint
// staff 建连时按 DB 实时 Actor 判定,决定是否加入管理团队房间
staff bool
// validate 应用层心跳时复查 token_version / 封禁状态,失效则立即断连
validate func() bool
// onHeartbeat 心跳顺带刷新 last_seen_at(SQL 60s 限频),
// 让只挂着 WS、不发 HTTP 请求的标签页也保持在线统计准确
onHeartbeat func()
// authorizeRoom 客户端动态订阅房间前的业务鉴权(如校验群成员身份);
// 仅 "chat:" 前缀房间允许订阅,user:/staff 等内部房间禁止
authorizeRoom func(room string) bool
send chan []byte
rooms map[string]struct{}
closed atomic.Bool
closeCh chan struct{}
}
func NewClient(
hub *Hub,
conn *websocket.Conn,
userID uint,
staff bool,
validate func() bool,
onHeartbeat func(),
authorizeRoom func(room string) bool,
) *Client {
return &Client{
hub: hub,
conn: conn,
userID: userID,
staff: staff,
validate: validate,
onHeartbeat: onHeartbeat,
authorizeRoom: authorizeRoom,
send: make(chan []byte, sendBufferSize),
rooms: make(map[string]struct{}),
closeCh: make(chan struct{}),
}
}
// Serve 注册到 Hub 后阻塞运行读写泵,断开时注销(调用方在升级成功后调用一次)。
func (c *Client) Serve(sendHello func(*Client)) {
cameOnline := c.hub.register(c)
sendHello(c)
if cameOnline {
c.hub.BroadcastPresence(c.userID, true)
}
go c.writePump()
c.readPump() // 阻塞至断连
wentOffline := c.hub.unregister(c)
c.forceClose()
if wentOffline {
c.hub.BroadcastPresence(c.userID, false)
}
}
// forceClose 幂等关闭连接并触发 closeCh
func (c *Client) forceClose() {
if c.closed.CompareAndSwap(false, true) {
_ = c.conn.Close()
close(c.closeCh)
}
}
// clientInbound 浏览器→服务端帧:ping 心跳;join/leave 动态订阅/退订房间
type clientInbound struct {
Type string `json:"type"`
Room string `json:"room,omitempty"`
}
// readPump 单读泵:Pong 续期 + 应用层心跳鉴权复查 + 房间订阅
func (c *Client) readPump() {
c.conn.SetReadLimit(maxMessageSize)
_ = c.conn.SetReadDeadline(time.Now().Add(pongWait))
c.conn.SetPongHandler(func(string) error {
return c.conn.SetReadDeadline(time.Now().Add(pongWait))
})
for {
_, raw, err := c.conn.ReadMessage()
if err != nil {
return
}
var msg clientInbound
if err := json.Unmarshal(raw, &msg); err != nil {
continue // 非法帧忽略,不踢连接
}
switch msg.Type {
case "ping":
// 心跳鉴权复查:被封禁/改密/降级强制下线后 ~25s 内断开
if c.validate != nil && !c.validate() {
_ = c.conn.WriteControl(
websocket.CloseMessage,
websocket.FormatCloseMessage(CloseAuthInvalid, "auth invalid"),
time.Now().Add(writeWait),
)
return
}
if c.onHeartbeat != nil {
c.onHeartbeat()
}
c.enqueueMustMarshal(Envelope{Type: EventPong, Data: map[string]int64{
"ts": time.Now().Unix(),
}})
case "join", "leave":
c.handleRoomFrame(msg)
}
}
}
// handleRoomFrame 处理房间订阅:只允许 chat:{id},且实时校验调用方业务身份
func (c *Client) handleRoomFrame(msg clientInbound) {
room := msg.Room
if room == "" || len(room) > 32 || !strings.HasPrefix(room, "chat:") {
return
}
idStr := strings.TrimPrefix(room, "chat:")
roomID, err := strconv.ParseUint(idStr, 10, 64)
if err != nil || roomID == 0 {
return
}
if c.authorizeRoom != nil && !c.authorizeRoom(room) {
return
}
if msg.Type == "join" {
c.hub.JoinClientRoom(c, room)
} else {
c.hub.LeaveClientRoom(c, room)
}
}
// writePump 单写泵:唯一持有连接写权限;协议层 Ping + 业务帧都从这里发出
func (c *Client) writePump() {
ticker := time.NewTicker(pingPeriod)
defer ticker.Stop()
for {
select {
case payload := <-c.send:
_ = c.conn.SetWriteDeadline(time.Now().Add(writeWait))
if err := c.conn.WriteMessage(websocket.TextMessage, payload); err != nil {
return
}
case <-ticker.C:
_ = c.conn.SetWriteDeadline(time.Now().Add(writeWait))
if err := c.conn.WriteMessage(websocket.PingMessage, nil); err != nil {
return
}
case <-c.closeCh:
return
}
}
}
// Send 业务帧入队(由唯一 writePump 发出);连接已关闭时静默丢弃
func (c *Client) Send(env Envelope) {
c.enqueueMustMarshal(env)
}
func (c *Client) enqueueMustMarshal(env Envelope) {
payload, err := json.Marshal(env)
if err != nil {
return
}
select {
case c.send <- payload:
case <-c.closeCh:
}
}

233
backend/realtime/hub.go Normal file
View File

@@ -0,0 +1,233 @@
// 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 校准角标)
)
// 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()
}
}
}