// 全站实时总线客户端(一期:在线状态 + 站点设置变更;二期:群聊 + 站内通知;三期:帖子流有新内容)。 // // - 单条 WebSocket(浏览器握手自动携带 HttpOnly cookie 鉴权),后续群聊/通知/ // 审核角标都复用本连接,按 message type 路由。 // - 连接状态机 + 指数退避重连;应用层心跳(协议层 Ping 由浏览器自动应答)。 // - 鉴权失效关闭码 4401 不自动重连(封禁/强制下线由 AccountGuard 统一处理)。 // - 在线集合:hello 帧给 staff 全量快照,之后 presence:update 增量维护。 // - 群聊:进入群详情页调用 joinRoom(roomID) 订阅 chat:{id} 房间, // 重连后会自动重订当前已订阅集合;离开页面/退出群调用 leaveRoom。 // - HTTP 仅作回前台 / 进页校准;实时靠 WS(通知、待审、主题、群聊未读)。 export const RT_HELLO = "hello"; export const RT_PONG = "pong"; export const RT_SETTINGS_CHANGED = "settings:changed"; export const RT_PRESENCE_UPDATE = "presence:update"; export const RT_CHAT_MESSAGE = "chat:message"; export const RT_CHAT_MEMBERSHIP = "chat:membership"; export const RT_CHAT_RECALLED = "chat:message_recalled"; export const RT_CHAT_UNREAD = "chat:unread"; export const RT_NOTIFICATION_NEW = "notification:new"; export const RT_FEED_CHANGED = "feed:changed"; export const RT_MODERATION_CHANGED = "moderation:changed"; export interface RtEnvelope { type: string; data?: T; } export interface RtHelloData { user_id: number; ts: number; // 仅管理团队成员会收到 online_user_ids?: number[]; } export interface RtPresenceData { user_id: number; online: boolean; } export interface RtSettingsData { accent: string; trust_reviewed_publish?: boolean; site_name?: string; site_description?: string; allow_register?: boolean; post_cooldown_hours?: number; code_block_auto_fold?: boolean; code_block_fold_lines?: number; } // 群聊新消息:data 为落库的 ChatMessage(含 sender) export interface RtChatMessageData { id: number; room_id: number; sender_id: number; content: string; reply_to_id?: number; reply_snap?: string; created_at: string; sender: { id: number; username: string; nickname: string; avatar: string; }; } // 群成员关系变更:被踢/群解散/被邀请 export interface RtChatMembershipData { room_id: number; status: "kicked" | "dissolved" | "invited"; } // 消息撤回 export interface RtChatRecalledData { id: number; room_id: number; sender_id?: number; recalled_by: number; recalled_at?: string; recaller?: { id: number; username: string; nickname: string; avatar: string; role?: string; }; } // 未读增量(推送给未在房内订阅的成员,用于 Header / 会话列表角标) export interface RtChatUnreadData { room_id: number; message_id: number; sender_id: number; preview?: string; delta: number; } // 新站内通知(如群聊 @):收到后应刷新通知列表/红点 export interface RtNotificationData { room_id?: number; message_id?: number; } // 帖子流有公开新内容(三期):首页 tip「有新内容」 export interface RtFeedChangedData { kind: "post" | "comment"; post_id: number; board_id: number; actor_id: number; } type RtHandler = (data: T) => void; type RtState = "idle" | "connecting" | "open" | "closed"; const HEARTBEAT_MS = 25_000; const RECONNECT_BASE_MS = 1_000; const RECONNECT_MAX_MS = 15_000; const CLOSE_AUTH_INVALID = 4401; // 开发环境 Next dev 的 rewrite 不透传 WebSocket Upgrade,直连后端 3001; // Cookie 不按端口隔离,localhost:3000 的 j13_token 会随握手发到 :3001。 // 生产环境前后端同源(反向代理),走同源 /api/ws;可用 NEXT_PUBLIC_WS_URL 覆盖。 function wsEndpoint(): string { const explicit = process.env.NEXT_PUBLIC_WS_URL; if (explicit) return explicit; if (typeof window === "undefined") return ""; const scheme = window.location.protocol === "https:" ? "wss://" : "ws://"; if (process.env.NODE_ENV !== "production") { return `${scheme}${window.location.hostname}:3001/api/ws`; } return `${scheme}${window.location.host}/api/ws`; } class RealtimeClient { private ws: WebSocket | null = null; private state: RtState = "idle"; // 是否应保持连接(登录态);登出或组件卸载时置 false 停止重连 private wanted = false; private retryCount = 0; private reconnectTimer: ReturnType | null = null; private heartbeatTimer: ReturnType | null = null; private readonly handlers = new Map>(); private readonly presenceHandlers = new Set>(); private onlineIds = new Set(); // 已订阅的群聊房间(roomID 集合):重连后会自动重订,断连期间不丢订阅意图 private readonly joinedRooms = new Set(); /** 建立连接(幂等);登录后调用 */ connect(): void { if (typeof window === "undefined") return; this.wanted = true; if (this.state === "open" || this.state === "connecting") return; this.clearReconnectTimer(); this.open(); window.addEventListener("online", this.handleNetworkBack); document.addEventListener("visibilitychange", this.handleVisible); } /** 主动断开(登出/卸载);不再自动重连 */ disconnect(): void { this.wanted = false; window.removeEventListener("online", this.handleNetworkBack); document.removeEventListener("visibilitychange", this.handleVisible); this.clearReconnectTimer(); this.clearHeartbeat(); if (this.ws) { this.ws.onclose = null; this.ws.close(1000, "client logout"); this.ws = null; } this.state = "idle"; // 登出时清空房间订阅,避免下次同标签页重新登录后误订旧群 this.joinedRooms.clear(); } isOpen(): boolean { return this.state === "open"; } /** 订阅指定事件,返回取消订阅函数 */ on(type: string, handler: RtHandler): () => void { let set = this.handlers.get(type); if (!set) { set = new Set(); this.handlers.set(type, set); } set.add(handler as RtHandler); return () => set!.delete(handler as RtHandler); } /** 订阅在线状态增量(用户 0↔1 连接跳变);快照请在订阅时读 getOnlineIds */ onPresence(handler: RtHandler): () => void { this.presenceHandlers.add(handler); return () => this.presenceHandlers.delete(handler); } /** 当前在线用户 ID 快照 */ getOnlineIds(): ReadonlySet { return this.onlineIds; } /** 订阅群聊房间(进入群详情页调用)。幂等;连接已 open 时立即发 join 帧。 * 后端会按 chat:{id} 实时校验成员身份,被踢后再 join 也会被拒绝。 */ joinRoom(roomId: number): void { if (!Number.isInteger(roomId) || roomId <= 0) return; this.joinedRooms.add(roomId); this.sendRoomFrame("join", roomId); } /** 退订群聊房间(离开群详情页/退出群调用)。幂等。 */ leaveRoom(roomId: number): void { if (!this.joinedRooms.delete(roomId)) return; this.sendRoomFrame("leave", roomId); } /** 退订全部群聊房间(切账号/被踢后清理) */ clearRooms(): void { for (const id of this.joinedRooms) { this.sendRoomFrame("leave", id); } this.joinedRooms.clear(); } private sendRoomFrame(type: "join" | "leave", roomId: number): void { if (this.state !== "open" || !this.ws) return; this.ws.send(JSON.stringify({ type, room: `chat:${roomId}` })); } private handleNetworkBack = (): void => { if (this.wanted && this.state === "closed") this.open(); }; private handleVisible = (): void => { if (!this.wanted || document.visibilityState !== "visible") return; // 后台标签页定时器可能被节流;重新可见时若连接已断则立即重连 if (this.state === "closed" || this.state === "idle") this.open(); }; private open(): void { if (typeof WebSocket === "undefined") return; this.state = "connecting"; let ws: WebSocket; try { ws = new WebSocket(wsEndpoint()); } catch { this.scheduleReconnect(); return; } this.ws = ws; ws.onopen = () => { this.state = "open"; this.retryCount = 0; this.clearHeartbeat(); this.heartbeatTimer = setInterval(() => { this.send({ type: "ping" }); }, HEARTBEAT_MS); // 重连后自动重订群聊房间(避免断线期间错过消息) // 后端会按 chat:{id} 实时校验成员身份,已退出/被踢的群会被静默拒绝 for (const roomId of this.joinedRooms) { this.sendRoomFrame("join", roomId); } }; ws.onmessage = (ev) => this.handleMessage(ev.data); ws.onerror = () => { // close 事件会紧随其后,重连逻辑统一在 onclose 处理 }; ws.onclose = (ev) => { this.clearHeartbeat(); if (this.ws === ws) this.ws = null; this.state = "closed"; // 鉴权失效:账号封禁/角色变更导致 token_version 失效, // 不重连(强制下线/重新登录由现有 HTTP 拦截链路处理) if (!this.wanted || ev.code === CLOSE_AUTH_INVALID) { this.wanted = false; return; } this.scheduleReconnect(); }; } private scheduleReconnect(): void { if (!this.wanted || this.reconnectTimer) return; const delay = Math.min( RECONNECT_MAX_MS, RECONNECT_BASE_MS * 2 ** this.retryCount ); // ±20% 抖动,避免后端重启时所有客户端同时重连 const jitter = delay * 0.2 * (Math.random() * 2 - 1); this.retryCount += 1; this.reconnectTimer = setTimeout(() => { this.reconnectTimer = null; if (this.wanted) this.open(); }, Math.max(500, delay + jitter)); } private clearReconnectTimer(): void { if (this.reconnectTimer) { clearTimeout(this.reconnectTimer); this.reconnectTimer = null; } } private clearHeartbeat(): void { if (this.heartbeatTimer) { clearInterval(this.heartbeatTimer); this.heartbeatTimer = null; } } private send(msg: RtEnvelope): void { if (this.state === "open" && this.ws) { this.ws.send(JSON.stringify(msg)); } } private handleMessage(raw: unknown): void { if (typeof raw !== "string") return; let msg: RtEnvelope; try { msg = JSON.parse(raw) as RtEnvelope; } catch { return; } if (!msg || typeof msg.type !== "string") return; // 内部状态事件先维护在线集合 if (msg.type === RT_HELLO) { const data = msg.data as RtHelloData | undefined; const ids = data?.online_user_ids; if (Array.isArray(ids)) { this.onlineIds = new Set(ids.filter((n) => Number.isInteger(n))); } } else if (msg.type === RT_PRESENCE_UPDATE) { const data = msg.data as RtPresenceData | undefined; if (data && typeof data.user_id === "number") { if (data.online) this.onlineIds.add(data.user_id); else this.onlineIds.delete(data.user_id); this.presenceHandlers.forEach((fn) => fn(data)); } } const set = this.handlers.get(msg.type); if (set) set.forEach((fn) => fn(msg.data)); } } // 全站单例(模块级,持久 root layout 下与标签页生命周期一致) export const realtime = new RealtimeClient();