import { ConnectedSocket, MessageBody, type OnGatewayConnection, type OnGatewayDisconnect, type OnGatewayInit, SubscribeMessage, WebSocketGateway, WebSocketServer, WsException, } from '@nestjs/websockets'; import { UseFilters, UseGuards } from '@nestjs/common'; import { MeetingRedisService } from './meeting-redis.service'; import { MeetingService } from '../meeting/meeting.service'; import { MeetingAuthGuard } from './meeting-auth.guard'; import { WsExceptionFilter } from '@/common/filters/ws-exception.filter'; import { LoggerService } from '@/plugins/logger/logger.service'; import type { ClientToServerMessageType, ErrorMessage, JoinRoomData, KickUserData, MeetingNamespace, MeetingRemoteSocket, MeetingSocket, MuteUserData, SucceedMessage, } from './types'; @WebSocketGateway({ // namespace 对应前端连接的 /meeting namespace: 'meeting', // 跨域(当前项目允许任意 origin) cors: { origin: '*' }, }) // 对所有 @SubscribeMessage 事件启用鉴权守卫 @UseGuards(MeetingAuthGuard) // WebSocket 异常统一格式化输出 @UseFilters(WsExceptionFilter) export class MeetingWebSocketGateway implements OnGatewayConnection, OnGatewayDisconnect, OnGatewayInit { @WebSocketServer() // 注入 Socket.IO namespace(带泛型,确保 server.in().fetchSockets() 等返回强类型) public server: MeetingNamespace | null = null; public constructor( // 房间/用户状态:Redis 持久化 private readonly redisService: MeetingRedisService, // shortUid/token 等业务能力 private readonly meetingService: MeetingService, // WebSocket 鉴权能力(用于 afterInit 中的 middleware) private readonly wsAuthGuard: MeetingAuthGuard, // 统一日志服务 private readonly logger: LoggerService ) {} private log(message: string): void { // 统一打到 MeetingWebSocket tag,便于检索 this.logger.info({}, message, 'MeetingWebSocket'); } /** * 查找目标用户的所有 Socket(本节点 + 跨节点) * @param roomId - 服务端房间 ID(Socket.IO 房间名) * @param targetUid - 目标用户声网 shortUid * @returns 目标用户的所有 Socket 数组 */ private async findTargetSockets(roomId: string, targetUid: number): Promise { // namespace(@WebSocketServer 注入)可能在启动早期为空,保护性返回 const ns = this.server; if (!ns) { return []; } // fetchSockets() 会返回本节点 Socket 或跨节点 RemoteSocket(包含 socket.data) const sockets = await ns.in(roomId).fetchSockets(); // 通过 socket.data.user.shortUid 精准定位目标用户(同一用户可能多端在线) return sockets.filter((s) => s.data.user?.shortUid === targetUid); } /** * 清理房间状态(如果已空) * 使用 Socket.IO 的 fetchSockets() 检查房间内是否有连接 * 注意:fetchSockets() 会返回本节点和跨节点的 socket * @param roomId - 服务端房间 ID(Socket.IO 房间名) */ private async cleanupRoomIfEmpty(roomId: string): Promise { const ns = this.server; if (!ns) { return; } // 使用 fetchSockets() 获取房间内的所有 socket(包括跨节点) const sockets = await ns.in(roomId).fetchSockets(); // 如果房间内没有 socket 连接,则清理 Redis 数据 if (sockets.length === 0) { await this.redisService.clearRoomAll(roomId); this.log(`[会议] 房间已空,已清理 Redis 状态:roomId=${roomId}`); } } /** * 初始化后,在 namespace 层添加连接鉴权 middleware * @param server - 注入的 Socket.IO namespace(带泛型,确保 server.in().fetchSockets() 等返回强类型) */ public afterInit(server: MeetingNamespace): void { // middleware:在 namespace 层做连接鉴权(用于 fetchSockets 时也能拿到 data.user) server.use(async (socket, next) => { try { // 解析握手信息并校验 Token,返回 user 信息 const user = await this.wsAuthGuard.validateToken(socket); // 写入 socket.data(会被 fetchSockets() 带回) socket.data.user = user; // 放行连接 next(); } catch (e) { console.error('鉴权失败====', e); // 交给 socket.io 处理为 connect_error,前端可据此处理 Unauthorized next(new WsException('Unauthorized')); } }); } public async handleConnection(socket: MeetingSocket): Promise { // 读取鉴权 middleware 写入的 user const user = socket.data.user; if (!user) { // 理论上不应该发生(有 guard + middleware),但仍兜底断开 this.logger.warn({}, `[MeetingWebSocket] WebSocket 认证失败,拒绝连接:${socket.id}`, 'MeetingWebSocket'); socket.disconnect(true); return; } // 连接建立日志(便于排查 shortUid/role 等) this.log(`新 WebSocket 连接建立:${socket.id} (userId=${user.userId}, shortUid=${user.shortUid}, role=${user.role})`); } public async handleDisconnect(socket: MeetingSocket): Promise { // Socket.IO 会自动将 socket 从房间移除,这里仅做日志 this.log(`WebSocket 连接断开:${socket.id}`); const roomId = socket.data.roomId; const shortUid = socket.data.user?.shortUid ? Number(socket.data.user.shortUid) : 0; const isHost = socket.data.user?.role !== 0; // 移除 socketId 映射(避免下次加入时误判为设备冲突) if (roomId && Number.isFinite(shortUid) && shortUid > 0) { await this.redisService.removeSocket(roomId, shortUid, socket.id).catch(() => {}); } // 如果是老师(创建者)主动断开连接,发送下课消息给所有人 if (roomId && isHost) { this.log(`老师(创建者)断开连接,发送下课消息:roomId=${roomId}`); // 设置课堂状态为已结束 await this.redisService.setClassStatus(roomId, 'finished'); // 通知所有人下课 this.server?.to(roomId).emit('message', { type: 'sev_class_ended', data: { fromRoomId: roomId } }); // 清理 Redis 数据 await this.redisService.clearRoomAll(roomId); } if (roomId) { this.cleanupRoomIfEmpty(roomId); } } @SubscribeMessage('client_join_room') public async handleJoinRoom(@MessageBody() data: JoinRoomData, @ConnectedSocket() socket: MeetingSocket): Promise { // 基础参数校验 if (!data?.courseRoomId) { this.logger.warn({}, '[会议] join_room 数据不完整', 'MeetingWebSocket'); // 下发标准错误包(前端统一处理) socket.emit('message', { type: 'error', data: { reason: '数据不完整' } } satisfies ErrorMessage); return; } // 统一为 string,避免 Redis key 不一致 const courseRoomId = String(data.courseRoomId); // 直接使用 courseRoomId 作为 roomId const roomId = courseRoomId; // 从鉴权后的 user 中获取 shortUid(鉴权通过就有短 ID) const shortUid = socket.data.user?.shortUid; const userName = socket.data.user?.userName || data.userName || '用户'; if (!shortUid) { this.logger.warn({}, '[会议] 用户未认证,无 shortUid', 'MeetingWebSocket'); socket.emit('message', { type: 'error', data: { reason: '认证失败' } } satisfies ErrorMessage); return; } // 黑名单校验(用短 UID 判断;被踢后不允许再次进入) const isKicked = await this.redisService.isUserKicked(roomId, shortUid); if (isKicked) { this.logger.warn({}, `[会议] 用户 ${shortUid} 已被踢出房间 ${roomId},拒绝重新加入`, 'MeetingWebSocket'); socket.emit('message', { type: 'sev_kick_user', data: { fromRoomId: roomId, reason: '您已被创建者移出会议,无法重新加入' }, } satisfies SucceedMessage); return; } // 检查该用户是否已在其他设备加入了房间 const existingSocketIds = await this.redisService.getSocketIds(roomId, shortUid); if (existingSocketIds.length > 0) { // 过滤掉当前 socket 自己的连接(同一设备刷新等情况) const otherDeviceSocketIds = existingSocketIds.filter((id) => id !== socket.id); if (otherDeviceSocketIds.length > 0) { // 验证旧连接是否真的存在(可能用户已断开但 Redis 未清理) // 使用 Promise.all 并行验证所有旧连接 const validationResults = await Promise.all( otherDeviceSocketIds.map(async (oldSocketId) => { const sockets = await this.server?.in(oldSocketId).fetchSockets(); return { oldSocketId, sockets, isValid: sockets ? sockets.length > 0 : false }; }) ); // 处理有效的旧连接:发送通知并断开 const validSocketIds: string[] = []; for (const result of validationResults) { if (result.isValid && result.sockets) { this.server?.to(result.oldSocketId).emit('message', { type: 'sev_device_conflict', data: { fromRoomId: roomId, reason: '您已在其他设备进入课程', targetUid: shortUid }, } satisfies SucceedMessage); result.sockets.forEach((s) => s.disconnect(true)); validSocketIds.push(result.oldSocketId); } } // 清理所有旧的 socketId(无论连接是否还存在) await Promise.all(otherDeviceSocketIds.map((oldSocketId) => this.redisService.removeSocket(roomId, shortUid, oldSocketId).catch(() => {}))); if (validSocketIds.length > 0) { this.log(`用户 ${userName}(短 UID:${shortUid}) 在其他设备加入,已踢出旧连接,房间 ${roomId}`); } } else { // 当前 socket 已在 Redis 中存在(同一设备刷新等情况),清理旧的并允许加入 this.log(`用户 ${userName}(短 UID:${shortUid}) 同一设备重新加入,房间 ${roomId}`); } } // 恢复用户状态(如果之前被禁麦/禁视频) const userState = await this.redisService.getUserState(roomId, shortUid); try { // 加入 Socket.IO 房间(用于广播/按房间查找) socket.join(roomId); socket.data.roomId = roomId; socket.data.courseRoomId = courseRoomId; // 注册 socketId 到 Redis(用于踢人、禁麦等控制功能) await this.redisService.addSocket(roomId, shortUid, socket.id); // 添加用户到房间用户列表 await this.redisService.addUserToRoom(roomId, shortUid); // 获取房间状态 const roomState = await this.redisService.getRoomState(roomId); const existedTeacherUid = Number.isFinite(roomState.teacherUid) ? Number(roomState.teacherUid) : 0; let teacherUid = existedTeacherUid > 0 ? existedTeacherUid : 0; if (teacherUid <= 0) { const homeworkId = Number(data?.homeworkId); const resolvedTeacherUid = Number.isFinite(homeworkId) && homeworkId > 0 ? await this.meetingService.getTeacherShortUidByHomeworkId(homeworkId) : null; teacherUid = resolvedTeacherUid && resolvedTeacherUid > 0 ? resolvedTeacherUid : shortUid; await this.redisService.setTeacherUid(roomId, teacherUid).catch(() => {}); } const existedSpeakerUid = Number.isFinite(roomState.speakerUid) ? Number(roomState.speakerUid) : 0; const speakerUid = existedSpeakerUid > 0 ? existedSpeakerUid : teacherUid; if (existedSpeakerUid <= 0) { await this.redisService.setSpeaker(roomId, speakerUid).catch(() => {}); } // 下发 join 成功包:包含 roomId/shortUid/tokenInfo/恢复状态/房间状态 socket.emit('message', { type: 'sev_join_room', data: { roomId, shortUid, tokenInfo: this.meetingService.generateToken(roomId, shortUid)!, isAudioMuted: userState.isAudioMuted, isVideoMuted: userState.isVideoMuted, classStatus: roomState.classStatus, screenShareUid: roomState.screenShareUid, speakerUid, teacherUid, }, }); // 成功日志 this.log(`用户 ${userName}(短 UID:${shortUid}) 加入房间 ${roomId}`); } catch (error) { // 记录错误后抛出,交由 WsExceptionFilter 统一处理 this.logger.error({}, `[会议] 为用户 ${shortUid} 分配短 UID 失败`, 'MeetingWebSocket'); throw error; } } @SubscribeMessage('client_leave_room') public async handleLeaveRoom(@MessageBody() data: { courseRoomId?: string }, @ConnectedSocket() socket: MeetingSocket): Promise { // Socket.IO 会自动处理离房逻辑,这里仅记录 this.log(`用户离开房间:${socket.id}`); const roomId: string | undefined = socket.data.roomId; const shortUid = socket.data.user?.shortUid ? Number(socket.data.user.shortUid) : 0; if (roomId && Number.isFinite(shortUid) && shortUid > 0) { // 移除 socketId 映射 await this.redisService.removeSocket(roomId, shortUid, socket.id).catch(() => {}); // 从房间用户列表移除 await this.redisService.removeUserFromRoom(roomId, shortUid).catch(() => {}); // 清理用户状态 await this.redisService.clearUserState(roomId, shortUid).catch(() => {}); } try { roomId && socket.leave(roomId); } catch {} socket.data.roomId = undefined; socket.data.courseRoomId = undefined; if (roomId) { await this.cleanupRoomIfEmpty(roomId); } } @SubscribeMessage('client_kick_user') public async handleKickUser(@MessageBody() data: KickUserData, @ConnectedSocket() socket: MeetingSocket): Promise { // 基础参数校验 if (!data?.targetUid || !data.roomId) { this.logger.warn({}, '[会议] kick_user 数据不完整', 'MeetingWebSocket'); return; } // 权限校验:仅创建者/管理员可踢人 const isHost = socket.data.user?.role !== 0; if (!isHost) { this.logger.warn({}, `[会议] 非创建者尝试踢人:userId=${socket.data.user?.userId}`, 'MeetingWebSocket'); return; } // 根据 Redis 中的 socketId 列表定位连接 const socketIds = await this.redisService.getSocketIds(data.roomId, data.targetUid); if (socketIds.length === 0) { this.logger.warn({}, `[会议] 未找到目标用户 uid: ${data.targetUid}`, 'MeetingWebSocket'); socket.emit('message', { type: 'error', data: { reason: '目标用户已离开房间或不存在' } } satisfies ErrorMessage); return; } // 获取用户名用于黑名单 const targetSockets = await this.findTargetSockets(data.roomId, data.targetUid); const targetUserName = targetSockets[0]?.data?.user?.userName || `用户-${data.targetUid}`; // 对该用户的所有连接发送踢出通知并断开连接 for (const id of socketIds) { this.server?.to(id).emit('message', { type: 'sev_kick_user', data: { fromRoomId: data.roomId, targetUid: data.targetUid } }); this.server?.in(id).disconnectSockets(true); } // 写入黑名单(用短 UID),禁止重连 await this.redisService.addToBlacklist(data.roomId, data.targetUid, targetUserName); // 成功日志 this.log(`用户 shortUid=${data.targetUid} 已被踢出房间 ${data.roomId}`); } @SubscribeMessage('client_mute_audio') public async handleMuteAudio(@MessageBody() data: MuteUserData, @ConnectedSocket() socket: MeetingSocket): Promise { // 基础参数校验 if (!data?.targetUid || !data.roomId) { this.logger.warn({}, '[会议] mute_audio 数据不完整', 'MeetingWebSocket'); return; } // 权限校验:仅创建者/管理员可禁麦 const isHost = socket.data.user?.role !== 0; if (!isHost) { this.logger.warn({}, `[会议] 非创建者尝试禁麦:userId=${socket.data.user?.userId}`, 'MeetingWebSocket'); return; } const socketIds = await this.redisService.getSocketIds(data.roomId, data.targetUid); if (socketIds.length === 0) { this.logger.warn({}, `[会议] 未找到目标用户 uid: ${data.targetUid}`, 'MeetingWebSocket'); socket.emit('message', { type: 'error', data: { reason: '目标用户已离开房间或不存在' } } satisfies ErrorMessage); return; } // 下发禁麦通知(由前端执行 Agora unpublish/setEnabled) for (const id of socketIds) { this.server?.to(id).emit('message', { type: 'sev_mute_audio', data: { fromRoomId: data.roomId, targetUid: data.targetUid } }); } // 写入 Redis:下次加入/重连时恢复禁麦状态 await this.redisService.setUserState(data.roomId, data.targetUid, { isAudioMuted: true }); this.log(`用户 ${data.targetUid} 已被禁麦,房间 ${data.roomId}`); } @SubscribeMessage('client_unmute_audio') public async handleUnmuteAudio(@MessageBody() data: MuteUserData, @ConnectedSocket() socket: MeetingSocket): Promise { // 基础参数校验 if (!data?.targetUid || !data.roomId) { this.logger.warn({}, '[会议] unmute_audio 数据不完整', 'MeetingWebSocket'); return; } // 权限校验:仅创建者/管理员可解除禁麦 const isHost = socket.data.user?.role !== 0; if (!isHost) { this.logger.warn({}, `[会议] 非创建者尝试解除禁麦:userId=${socket.data.user?.userId}`, 'MeetingWebSocket'); return; } const socketIds = await this.redisService.getSocketIds(data.roomId, data.targetUid); if (socketIds.length === 0) { this.logger.warn({}, `[会议] 未找到目标用户 uid: ${data.targetUid}`, 'MeetingWebSocket'); socket.emit('message', { type: 'error', data: { reason: '目标用户已离开房间或不存在' } } satisfies ErrorMessage); return; } // 下发解除禁麦通知 for (const id of socketIds) { this.server?.to(id).emit('message', { type: 'sev_unmute_audio', data: { fromRoomId: data.roomId, targetUid: data.targetUid } }); } // 写入 Redis:下次加入/重连时恢复状态 await this.redisService.setUserState(data.roomId, data.targetUid, { isAudioMuted: false }); this.log(`用户 ${data.targetUid} 已被解除禁麦,房间 ${data.roomId}`); } @SubscribeMessage('client_mute_video') public async handleMuteVideo(@MessageBody() data: MuteUserData, @ConnectedSocket() socket: MeetingSocket): Promise { // 基础参数校验 if (!data?.targetUid || !data.roomId) { this.logger.warn({}, '[会议] mute_video 数据不完整', 'MeetingWebSocket'); return; } // 权限校验:仅创建者/管理员可禁视频 const isHost = socket.data.user?.role !== 0; if (!isHost) { this.logger.warn({}, `[会议] 非创建者尝试禁视频:userId=${socket.data.user?.userId}`, 'MeetingWebSocket'); return; } const socketIds = await this.redisService.getSocketIds(data.roomId, data.targetUid); if (socketIds.length === 0) { this.logger.warn({}, `[会议] 未找到目标用户 uid: ${data.targetUid}`, 'MeetingWebSocket'); socket.emit('message', { type: 'error', data: { reason: '目标用户已离开房间或不存在' } } satisfies ErrorMessage); return; } // 下发禁视频通知(由前端执行 Agora unpublish/setEnabled) for (const id of socketIds) { this.server?.to(id).emit('message', { type: 'sev_mute_video', data: { fromRoomId: data.roomId, targetUid: data.targetUid } }); } // 写入 Redis:下次加入/重连时恢复禁视频状态 await this.redisService.setUserState(data.roomId, data.targetUid, { isVideoMuted: true }); this.log(`用户 ${data.targetUid} 已被禁视频,房间 ${data.roomId}`); } @SubscribeMessage('client_unmute_video') public async handleUnmuteVideo(@MessageBody() data: MuteUserData, @ConnectedSocket() socket: MeetingSocket): Promise { // 基础参数校验 if (!data?.targetUid || !data.roomId) { this.logger.warn({}, '[会议] unmute_video 数据不完整', 'MeetingWebSocket'); return; } // 权限校验:仅创建者/管理员可解除禁视频 const isHost = socket.data.user?.role !== 0; if (!isHost) { this.logger.warn({}, `[会议] 非创建者尝试解除禁视频:userId=${socket.data.user?.userId}`, 'MeetingWebSocket'); return; } const socketIds = await this.redisService.getSocketIds(data.roomId, data.targetUid); if (socketIds.length === 0) { this.logger.warn({}, `[会议] 未找到目标用户 uid: ${data.targetUid}`, 'MeetingWebSocket'); socket.emit('message', { type: 'error', data: { reason: '目标用户已离开房间或不存在' } } satisfies ErrorMessage); return; } // 下发解除禁视频通知 for (const id of socketIds) { this.server?.to(id).emit('message', { type: 'sev_unmute_video', data: { fromRoomId: data.roomId, targetUid: data.targetUid } }); } // 写入 Redis:下次加入/重连时恢复状态 await this.redisService.setUserState(data.roomId, data.targetUid, { isVideoMuted: false }); this.log(`用户 ${data.targetUid} 已被解除禁视频,房间 ${data.roomId}`); } /** * 处理设置全员主屏消息(仅创建者) */ @SubscribeMessage('client_set_main_video') public async handleSetMainVideo(@MessageBody() data: MuteUserData, @ConnectedSocket() socket: MeetingSocket): Promise { // 参数校验 if (!data?.targetUid || !data.roomId) { this.logger.warn({}, '[会议] set_main_video 数据不完整', 'MeetingWebSocket'); return; } // 权限校验 const isHost = socket.data.user?.role !== 0; if (!isHost) { this.logger.warn({}, `[会议] 非创建者尝试设置主屏:userId=${socket.data.user?.userId}`, 'MeetingWebSocket'); return; } // 广播给整个房间 const ns = this.server; if (!ns) { return; } ns.to(data.roomId).emit('message', { type: 'sev_set_main_video', data: { fromRoomId: data.roomId, targetUid: data.targetUid } }); await this.redisService.setSpeaker(data.roomId, Number(data.targetUid)).catch(() => {}); this.log(`已设置全员主屏:targetUid=${data.targetUid} 房间 ${data.roomId}`); } /** * 老师上课:允许推流 */ @SubscribeMessage('client_start_class') public async handleStartClass(@MessageBody() data: { roomId: string }, @ConnectedSocket() socket: MeetingSocket): Promise { if (!data?.roomId) { return; } const isHost = socket.data.user?.role !== 0; if (!isHost) { return; } await this.redisService.setClassStatus(data.roomId, 'in_class'); // 通知所有人可以开始推流了 this.server?.to(data.roomId).emit('message', { type: 'sev_class_started', data: { fromRoomId: data.roomId } }); this.log(`房间 ${data.roomId} 上课开始`); } /** * 老师下课:停止推流并清理本节课数据(保持 socket 连接) */ @SubscribeMessage('client_end_class') public async handleEndClass(@MessageBody() data: { roomId: string }, @ConnectedSocket() socket: MeetingSocket): Promise { if (!data?.roomId) { return; } const isHost = socket.data.user?.role !== 0; if (!isHost) { return; } await this.redisService.setClassStatus(data.roomId, 'finished'); // 通知所有人结束推流(由前端停止/离开 RTC) this.server?.to(data.roomId).emit('message', { type: 'sev_class_ended', data: { fromRoomId: data.roomId } }); // 清理本节课全部 Redis 数据(不清理黑名单) await this.redisService.clearClassData(data.roomId); this.log(`房间 ${data.roomId} 下课并清理 Redis 数据`); } /** * 处理开始投屏消息(仅创建者) */ @SubscribeMessage('client_start_screen_share') public async handleStartScreenShare(@MessageBody() data: { roomId: string }, @ConnectedSocket() socket: MeetingSocket): Promise { if (!data?.roomId) { return; } const isHost = socket.data.user?.role !== 0; if (!isHost) { this.logger.warn({}, `[会议] 非创建者尝试开始投屏:userId=${socket.data.user?.userId}`, 'MeetingWebSocket'); return; } const targetUid = socket.data.user?.shortUid ?? 0; // 保存投屏状态到 Redis await this.redisService.setScreenSharing(data.roomId, targetUid); // 广播给整个房间:有人开始投屏 this.server?.to(data.roomId).emit('message', { type: 'sev_start_screen_share', data: { fromRoomId: data.roomId, targetUid } }); this.log(`创建者开始投屏:roomId=${data.roomId}, shortUid=${targetUid}`); } /** * 处理停止投屏消息(仅创建者) */ @SubscribeMessage('client_stop_screen_share') public async handleStopScreenShare(@MessageBody() data: { roomId: string }, @ConnectedSocket() socket: MeetingSocket): Promise { if (!data?.roomId) { return; } const isHost = socket.data.user?.role !== 0; if (!isHost) { this.logger.warn({}, `[会议] 非创建者尝试停止投屏:userId=${socket.data.user?.userId}`, 'MeetingWebSocket'); return; } const targetUid = socket.data.user?.shortUid ?? 0; // 清除投屏状态 await this.redisService.setScreenSharing(data.roomId, null); // 广播给整个房间:有人停止投屏 this.server?.to(data.roomId).emit('message', { type: 'sev_stop_screen_share', data: { fromRoomId: data.roomId, targetUid } }); this.log(`创建者停止投屏:roomId=${data.roomId}, shortUid=${targetUid}`); } }