Files
jiang13-bbs/frontend/lib/realtime.ts

409 lines
13 KiB
TypeScript
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// 全站实时总线客户端(一期:在线状态 + 站点设置变更;二期:群聊 + 站内通知;三期:帖子流有新内容)。
//
// - 单条 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<T = unknown> {
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;
site_wordmark?: string;
site_slogan?: string;
site_keywords?: string[];
logo_light_url?: string;
logo_dark_url?: string;
favicon_url?: string;
brand_mark?: string;
brand_logo_size?: string;
brand_logo_fit?: string;
footer_links?: { label: string; url: string; new_tab: boolean }[];
allow_register?: boolean;
allow_comments?: boolean;
allow_messages?: boolean;
post_cooldown_hours?: number;
code_block_auto_fold?: boolean;
code_block_fold_lines?: number;
ui_animations?: boolean;
anim_code_fold?: boolean;
anim_smooth_scroll?: boolean;
anim_chrome?: boolean;
post_link_new_tab?: boolean;
attachment_ext_limit?: boolean;
attachment_exts?: string[];
attachment_max_mb?: number;
attachment_max_count?: number;
image_max_mb?: number;
bg_site_url?: string;
bg_site_mode?: string;
bg_admin_url?: string;
bg_admin_mode?: string;
}
// 群聊新消息: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 / 会话列表角标)
// delta 通常为 +1(新消息);撤回未读消息时为 -1,此时不含 preview,不得当新消息处理
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<T = unknown> = (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<typeof setTimeout> | null = null;
private heartbeatTimer: ReturnType<typeof setInterval> | null = null;
// Strict Mode / HMR 会瞬时卸载→重挂;延迟断开,避免握手未完成就被 close 刷控制台
private releaseTimer: ReturnType<typeof setTimeout> | null = null;
private readonly handlers = new Map<string, Set<RtHandler>>();
private readonly presenceHandlers = new Set<RtHandler<RtPresenceData>>();
private onlineIds = new Set<number>();
// 已订阅的群聊房间(roomID 集合):重连后会自动重订,断连期间不丢订阅意图
private readonly joinedRooms = new Set<number>();
/** 建立连接(幂等);登录后调用 */
connect(): void {
if (typeof window === "undefined") return;
this.clearReleaseTimer();
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);
}
/**
* 组件卸载时调用:短暂延迟后再真正断开。
* 若 Strict Mode / HMR 随即重新 connect,会取消本次断开,保留握手中的连接。
*/
release(delayMs = 120): void {
this.clearReleaseTimer();
this.releaseTimer = setTimeout(() => {
this.releaseTimer = null;
this.disconnect();
}, delayMs);
}
/** 主动断开(登出/卸载);不再自动重连 */
disconnect(): void {
this.clearReleaseTimer();
this.wanted = false;
window.removeEventListener("online", this.handleNetworkBack);
document.removeEventListener("visibilitychange", this.handleVisible);
this.clearReconnectTimer();
this.clearHeartbeat();
if (this.ws) {
this.ws.onclose = null;
// CONNECTING 时直接 close 会触发浏览器警告;改等 open 后再关
if (this.ws.readyState === WebSocket.CONNECTING) {
const pending = this.ws;
pending.onopen = () => {
pending.close(1000, "client logout");
};
} else {
this.ws.close(1000, "client logout");
}
this.ws = null;
}
this.state = "idle";
// 登出时清空房间订阅,避免下次同标签页重新登录后误订旧群
this.joinedRooms.clear();
}
isOpen(): boolean {
return this.state === "open";
}
/** 订阅指定事件,返回取消订阅函数 */
on<T>(type: string, handler: RtHandler<T>): () => 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<RtPresenceData>): () => void {
this.presenceHandlers.add(handler);
return () => this.presenceHandlers.delete(handler);
}
/** 当前在线用户 ID 快照 */
getOnlineIds(): ReadonlySet<number> {
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 clearReleaseTimer(): void {
if (this.releaseTimer) {
clearTimeout(this.releaseTimer);
this.releaseTimer = 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();