import WebSocket from 'ws'; import { createHash } from 'crypto'; import { chatServerUrl, fetchChannelInfo, fetchChatAccessToken, fetchLiveChatInfo, } from './channel'; import { addChatMessage, addDiceFromDonation, clearChatFeed, getRoom, getRoomViewerCount, hasSeenDonationKey, markSeenDonationKey, notifyRoom, setRoomMeta, } from '../rooms/store'; const ChatCmd = { PING: 0, PONG: 10000, CONNECT: 100, CONNECTED: 10100, REQUEST_RECENT_CHAT: 5101, RECENT_CHAT: 15101, CHAT: 93101, DONATION: 93102, } as const; const ChatType = { TEXT: 1, DONATION: 10, SUBSCRIPTION: 11, } as const; type SessionHandle = { roomId: string; ws: WebSocket; chatChannelId: string; sid: string | null; pingTimer: ReturnType | null; pollTimer: ReturnType | null; closed: boolean; /** 이번 연결에서 이미 피드에 넣은 메시지 키 */ seenMessageKeys: Set; }; const SEEN_KEYS_MAX = 800; const RECENT_CHAT_COUNT = 50; const globalForSessions = globalThis as typeof globalThis & { __drinkmarbleChatSessions?: Map; }; function getSessions() { if (!globalForSessions.__drinkmarbleChatSessions) { globalForSessions.__drinkmarbleChatSessions = new Map(); } return globalForSessions.__drinkmarbleChatSessions; } function clearPing(handle: SessionHandle) { if (handle.pingTimer) clearTimeout(handle.pingTimer); handle.pingTimer = null; } function schedulePing(handle: SessionHandle) { clearPing(handle); handle.pingTimer = setTimeout(() => { if (handle.closed || handle.ws.readyState !== WebSocket.OPEN) return; handle.ws.send(JSON.stringify({ cmd: ChatCmd.PING, ver: '2' })); schedulePing(handle); }, 20000); } function parseMaybeJson(value: unknown): T | null { if (value == null) return null; if (typeof value === 'object') return value as T; if (typeof value !== 'string') return null; try { return JSON.parse(value) as T; } catch { return null; } } function messageTimeMs(chat: Record): number | null { const raw = chat.msgTime ?? chat.messageTime; const n = typeof raw === 'string' ? Number(raw) : Number(raw); if (!Number.isFinite(n) || n <= 0) return null; return n; } function messageKey( chat: Record, kind: 'donation' | 'chat', ): string | null { const id = chat.msgUid ?? chat.messageId ?? null; if (id != null && String(id).length > 0) { return hashMessageKey(`${kind}:${String(id)}`); } const time = messageTimeMs(chat); const profile = parseMaybeJson<{ nickname?: string; userIdHash?: string }>( chat.profile, ); const extras = parseMaybeJson<{ payAmount?: number }>(chat.extras); const who = profile?.userIdHash || profile?.nickname || ''; const msg = String(chat.msg ?? chat.content ?? ''); if (time == null && !who && !msg) return null; if (kind === 'donation') { return hashMessageKey( `donation:${time ?? 0}:${who}:${Number(extras?.payAmount ?? 0)}:${msg}`, ); } return hashMessageKey(`chat:${time ?? 0}:${who}:${msg}`); } /** 중복 검사용 키에는 사용자 해시·닉네임·메시지 원문을 보관하지 않는다. */ function hashMessageKey(value: string): string { return createHash('sha256').update(value).digest('base64url'); } function rememberMessageKey(handle: SessionHandle, key: string | null): boolean { if (!key) return true; if (handle.seenMessageKeys.has(key)) return false; handle.seenMessageKeys.add(key); if (handle.seenMessageKeys.size > SEEN_KEYS_MAX) { const oldest = handle.seenMessageKeys.values().next().value; if (oldest != null) handle.seenMessageKeys.delete(oldest); } return true; } function handleDonationBody( roomId: string, handle: SessionHandle, chat: Record, isRecent: boolean, ): boolean { const key = messageKey(chat, 'donation'); if (!rememberMessageKey(handle, key)) return false; const extras = parseMaybeJson<{ payAmount?: number; donationType?: string; }>(chat.extras); const profile = parseMaybeJson<{ nickname?: string }>(chat.profile); const payAmount = Number(extras?.payAmount ?? 0); if (!Number.isFinite(payAmount) || payAmount <= 0) return false; // 히스토리 후원: 피드에만 표시하고, 이미 본 것으로 표시해 재적용 방지 if (isRecent) { if (key) markSeenDonationKey(roomId, key); addDiceFromDonation(roomId, { payAmount, donatorNickname: profile?.nickname, donationText: String(chat.msg ?? chat.content ?? ''), applyDice: false, }); return true; } if (key && hasSeenDonationKey(roomId, key)) { return false; } if (key) markSeenDonationKey(roomId, key); addDiceFromDonation(roomId, { payAmount, donatorNickname: profile?.nickname, donationText: String(chat.msg ?? chat.content ?? ''), applyDice: true, }); return true; } function handleChatBody( roomId: string, handle: SessionHandle, chat: Record, ): boolean { const key = messageKey(chat, 'chat'); if (!rememberMessageKey(handle, key)) return false; const profile = parseMaybeJson<{ nickname?: string }>(chat.profile); const message = String(chat.msg ?? chat.content ?? '').trim(); if (!message) return false; addChatMessage(roomId, { nickname: profile?.nickname, message, }); return true; } function processMessageList( roomId: string, handle: SessionHandle, body: Record | unknown[] | undefined, cmd: number, ) { const isRecent = cmd === ChatCmd.RECENT_CHAT; const list = Array.isArray(body) ? body : Array.isArray((body as { messageList?: unknown[] })?.messageList) ? ((body as { messageList: unknown[] }).messageList) : body ? [body] : []; let changed = false; for (const item of list) { if (!item || typeof item !== 'object') continue; const chat = item as Record; const type = Number(chat.msgTypeCode ?? chat.messageTypeCode ?? 0); if (cmd === ChatCmd.DONATION || type === ChatType.DONATION) { if (handleDonationBody(roomId, handle, chat, isRecent)) changed = true; } else if ( type === ChatType.TEXT || type === ChatType.SUBSCRIPTION || !type ) { if (handleChatBody(roomId, handle, chat)) changed = true; } } // 최근 채팅 50개를 메시지마다 SSE로 보내면 주사위 Canvas가 깨질 수 있어 한 번만 알림 if (changed) notifyRoom(roomId); } function requestRecentChat(handle: SessionHandle) { if ( handle.closed || !handle.sid || handle.ws.readyState !== WebSocket.OPEN ) { return; } handle.ws.send( JSON.stringify({ ver: '2', cmd: ChatCmd.REQUEST_RECENT_CHAT, svcid: 'game', cid: handle.chatChannelId, sid: handle.sid, tid: 2, bdy: { recentMessageCount: RECENT_CHAT_COUNT }, }), ); } function onSocketMessage(roomId: string, handle: SessionHandle, raw: WebSocket.RawData) { let json: { cmd?: number; bdy?: unknown; }; try { json = JSON.parse(String(raw)); } catch { return; } const cmd = json.cmd; const body = json.bdy as Record | unknown[] | undefined; if (cmd === ChatCmd.CONNECTED) { const sid = body && !Array.isArray(body) && typeof body.sid === 'string' ? body.sid : null; handle.sid = sid; handle.seenMessageKeys.clear(); clearChatFeed(roomId); setRoomMeta(roomId, { chatStatus: 'connected', chatError: null }); notifyRoom(roomId); schedulePing(handle); requestRecentChat(handle); return; } if (cmd === ChatCmd.PING) { if (handle.ws.readyState === WebSocket.OPEN) { handle.ws.send(JSON.stringify({ cmd: ChatCmd.PONG, ver: '2' })); } return; } if ( cmd === ChatCmd.CHAT || cmd === ChatCmd.DONATION || cmd === ChatCmd.RECENT_CHAT ) { processMessageList(roomId, handle, body, cmd ?? 0); schedulePing(handle); } } async function openSocket(roomId: string, channelId: string): Promise { const live = await fetchLiveChatInfo(channelId); const accessToken = await fetchChatAccessToken(live.chatChannelId); const url = chatServerUrl(live.chatChannelId); setRoomMeta(roomId, { chatChannelId: live.chatChannelId, chatStatus: 'connecting', chatError: null, }); notifyRoom(roomId); const ws = new WebSocket(url, { headers: { Origin: 'https://chzzk.naver.com', 'User-Agent': 'Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/122.0.0.0 Safari/537.36', }, }); const handle: SessionHandle = { roomId, ws, chatChannelId: live.chatChannelId, sid: null, pingTimer: null, pollTimer: null, closed: false, seenMessageKeys: new Set(), }; await new Promise((resolve, reject) => { const timer = setTimeout(() => { reject(new Error('채팅 웹소켓 연결 시간 초과')); }, 12000); ws.on('open', () => { ws.send( JSON.stringify({ ver: '2', cmd: ChatCmd.CONNECT, svcid: 'game', cid: live.chatChannelId, tid: 1, bdy: { uid: null, devType: 2001, accTkn: accessToken, auth: 'READ', }, }), ); }); ws.on('message', (data) => { try { const json = JSON.parse(String(data)) as { cmd?: number }; if (json.cmd === ChatCmd.CONNECTED) { clearTimeout(timer); onSocketMessage(roomId, handle, data); resolve(); return; } } catch { // fallthrough } onSocketMessage(roomId, handle, data); }); ws.on('error', (err) => { clearTimeout(timer); reject(err instanceof Error ? err : new Error('웹소켓 오류')); }); ws.on('close', () => { handle.closed = true; clearPing(handle); if (handle.pollTimer) clearInterval(handle.pollTimer); const current = getSessions().get(roomId); if (current === handle) { getSessions().delete(roomId); setRoomMeta(roomId, { chatStatus: 'disconnected' }); notifyRoom(roomId); } }); }); // chatChannelId 변경 폴링 (방송 재시작 대응) handle.pollTimer = setInterval(() => { void (async () => { try { const next = await fetchLiveChatInfo(channelId); if (next.chatChannelId && next.chatChannelId !== handle.chatChannelId) { await disconnectRoomChat(roomId); await connectRoomChat(roomId); } } catch { // ignore poll errors } })(); }, 30000); getSessions().set(roomId, handle); return handle; } export function isRoomChatLive(roomId: string) { const handle = getSessions().get(roomId); return Boolean( handle && !handle.closed && handle.ws.readyState === WebSocket.OPEN && handle.sid, ); } const globalForDisconnect = globalThis as typeof globalThis & { __drinkmarbleChatDisconnectTimers?: Map>; }; function getDisconnectTimers() { if (!globalForDisconnect.__drinkmarbleChatDisconnectTimers) { globalForDisconnect.__drinkmarbleChatDisconnectTimers = new Map(); } return globalForDisconnect.__drinkmarbleChatDisconnectTimers; } /** 빠른 재입장 시 이전 leave가 새 연결을 끊지 않도록 예약 해제 */ export function cancelScheduledDisconnectRoomChat(roomId: string) { const timers = getDisconnectTimers(); const t = timers.get(roomId); if (!t) return; clearTimeout(t); timers.delete(roomId); } /** * 시청자가 0이 되어도 바로 끊지 않고 잠시 대기. * 그 사이 재입장하면 채팅 세션을 유지한다. */ export function scheduleDisconnectRoomChat(roomId: string, delayMs = 4000) { cancelScheduledDisconnectRoomChat(roomId); const timers = getDisconnectTimers(); const t = setTimeout(() => { timers.delete(roomId); if (getRoomViewerCount(roomId) > 0) return; void disconnectRoomChat(roomId); }, delayMs); timers.set(roomId, t); } /** 살아 있으면 유지(최근 채팅만 재요청), 아니면 재연결 */ export async function ensureRoomChat(roomId: string) { cancelScheduledDisconnectRoomChat(roomId); const room = getRoom(roomId); if (!room) throw new Error('방을 찾을 수 없습니다.'); if (!room.settings.chzzkEnabled) { throw new Error('이 방은 치지직 연동이 꺼져 있습니다.'); } if (isRoomChatLive(roomId)) { const handle = getSessions().get(roomId); if (handle) requestRecentChat(handle); setRoomMeta(roomId, { chatStatus: 'connected', chatError: null }); notifyRoom(roomId); return; } await connectRoomChat(roomId); } export async function connectRoomChat(roomId: string) { cancelScheduledDisconnectRoomChat(roomId); const room = getRoom(roomId); if (!room) throw new Error('방을 찾을 수 없습니다.'); if (!room.settings.chzzkEnabled) { throw new Error('이 방은 치지직 연동이 꺼져 있습니다.'); } const channelId = room.settings.channelId; if (!channelId) throw new Error('channelId가 없습니다.'); await disconnectRoomChat(roomId); try { const info = await fetchChannelInfo(channelId); setRoomMeta(roomId, { channelName: info.channelName, chatStatus: 'connecting' }); notifyRoom(roomId); await openSocket(roomId, channelId); } catch { // 외부 라이브러리 오류에는 URL·토큰·식별자가 섞일 수 있어 원문을 보관하지 않는다. const message = '채팅 연결에 실패했습니다. 방송 상태를 확인해 주세요.'; setRoomMeta(roomId, { chatStatus: 'error', chatError: message }); notifyRoom(roomId); throw new Error(message); } } export async function disconnectRoomChat(roomId: string) { cancelScheduledDisconnectRoomChat(roomId); const handle = getSessions().get(roomId); if (!handle) { setRoomMeta(roomId, { chatStatus: 'disconnected' }); notifyRoom(roomId); return; } handle.closed = true; clearPing(handle); if (handle.pollTimer) clearInterval(handle.pollTimer); getSessions().delete(roomId); try { handle.ws.close(); } catch { // ignore } setRoomMeta(roomId, { chatStatus: 'disconnected' }); notifyRoom(roomId); }