/** * mediasoup-client 封装层 * * 职责: * - 封装 Device/sendTransport/recvTransport/Producer/Consumer 的生命周期 * - 把 Transport 的 connect/produce 回调桥接到 WS 信令(sendWithAck) * - 对上层(store/meeting.js)暴露同步友好的异步 API * * 平台约束(Task 9 决策 Q1=a1_h5_only): * - mediasoup-client 仅在 H5 端可用;非 H5 平台调用时会在构造时抛 ERR_PLATFORM * - import 语句通过 uni-app 条件编译注释限定在 H5 构建 * * 使用方式(仅限 store/meeting.js): * import { createMediaEngine } from '@/utils/mediasoup-client' * const engine = markRaw(createMediaEngine({ roomCode, userId, sendWithAck })) * await engine.loadDevice(rtpCapabilities) * const sendTransport = await engine.ensureSendTransport() * const producer = await engine.produce({ kind: 'audio', track }) */ // #ifdef H5 import { Device } from 'mediasoup-client' // #endif import { MEETING_WS_TRANSPORT_CREATE, MEETING_WS_TRANSPORT_CONNECT, MEETING_WS_PRODUCE_START, MEETING_WS_CONSUME_START, MEETING_WS_CONSUME_RESUME, MEETING_WS_PRODUCER_CLOSE } from '@/constants/meeting' /** * 创建 MediaEngine 实例 * * @param {Object} options * @param {string} options.roomCode - 会议号 * @param {number} options.userId - 当前用户 ID * @param {Function} options.sendWithAck - WS 发送带 ACK 等待的方法(由 services/websocket.js 提供) * @param {Function} [options.logger] - 日志函数,默认 console.log * @returns {Object} MediaEngine 实例 */ export function createMediaEngine({ roomCode, userId, sendWithAck, logger = defaultLogger }) { if (!roomCode) throw new Error('[MediaEngine] roomCode 不能为空') if (!userId) throw new Error('[MediaEngine] userId 不能为空') if (typeof sendWithAck !== 'function') { throw new Error('[MediaEngine] sendWithAck 必须是函数') } // #ifndef H5 throw new Error('[MediaEngine] 当前平台不支持 mediasoup-client(仅支持 H5)') // #endif // #ifdef H5 /** mediasoup Device 实例;Device.load 后持有 Router rtpCapabilities */ let device = null /** 发送 Transport(本地 Producer 用),延迟创建,复用 */ let sendTransport = null /** 接收 Transport(本地 Consumer 用),延迟创建,复用 */ let recvTransport = null /** 本地 Producer 索引 Map */ const producers = new Map() /** 本地 Consumer 索引 Map */ const consumers = new Map() /** 是否已关闭,防止重复调用 close */ let closed = false const ensureNotClosed = () => { if (closed) { throw new Error('[MediaEngine] 已关闭') } } const ensureDeviceLoaded = () => { ensureNotClosed() if (!device || !device.loaded) { throw new Error('[MediaEngine] Device 尚未 load') } } /** * 加载 Device * @param {Object} routerRtpCapabilities - 后端 JoinRoom 响应中的 rtp_capabilities */ const loadDevice = async (routerRtpCapabilities) => { ensureNotClosed() if (device && device.loaded) { logger('debug', '[MediaEngine] Device 已 load,跳过') return } if (!routerRtpCapabilities) { throw new Error('[MediaEngine] routerRtpCapabilities 不能为空') } device = new Device() await device.load({ routerRtpCapabilities }) logger('info', '[MediaEngine] Device load 成功', { canProduceAudio: device.canProduce('audio'), canProduceVideo: device.canProduce('video') }) } /** 返回 Device.rtpCapabilities(用于后端 CreateConsumer) */ const getRtpCapabilities = () => { ensureDeviceLoaded() return device.rtpCapabilities } /** * 通过 WS 请求后端创建 Transport,返回 mediasoup-client 可用的 Transport 元信息 * @param {'send'|'recv'} direction */ const requestTransportInfo = async (direction) => { const info = await sendWithAck(MEETING_WS_TRANSPORT_CREATE, { room_code: roomCode, direction }) if (!info || !info.id) { throw new Error(`[MediaEngine] 创建 ${direction} Transport 失败:后端返回为空`) } return info } /** * 绑定 sendTransport 的 connect + produce 回调(把本地事件桥到 WS 信令) * mediasoup-client 约定: * - transport.on('connect', ({ dtlsParameters }, callback, errback)) * 需要在 DTLS 握手完成前告知远端 DTLS 参数 * - transport.on('produce', ({ kind, rtpParameters, appData }, callback, errback)) * 需要在 Producer 创建前把 RTP 参数发给远端,拿到远端分配的 producerId 回调 callback({ id }) */ const bindSendTransportEvents = (transport) => { transport.on('connect', ({ dtlsParameters }, callback, errback) => { sendWithAck(MEETING_WS_TRANSPORT_CONNECT, { room_code: roomCode, transport_id: transport.id, dtls_parameters: dtlsParameters }).then(() => callback()).catch((err) => { logger('error', '[MediaEngine] sendTransport connect 失败', err) errback(err) }) }) transport.on('produce', ({ kind, rtpParameters, appData }, callback, errback) => { sendWithAck(MEETING_WS_PRODUCE_START, { room_code: roomCode, transport_id: transport.id, kind, rtp_parameters: rtpParameters, app_data: appData || {} }).then((resp) => { if (!resp || !resp.producer_id) { errback(new Error('后端返回 producer_id 为空')) return } callback({ id: resp.producer_id }) }).catch((err) => { logger('error', '[MediaEngine] sendTransport produce 失败', err) errback(err) }) }) transport.on('connectionstatechange', (state) => { logger('debug', `[MediaEngine] sendTransport state=${state}`) }) } /** recvTransport 只需要桥接 connect 回调(consume 由 store 主动调起) */ const bindRecvTransportEvents = (transport) => { transport.on('connect', ({ dtlsParameters }, callback, errback) => { sendWithAck(MEETING_WS_TRANSPORT_CONNECT, { room_code: roomCode, transport_id: transport.id, dtls_parameters: dtlsParameters }).then(() => callback()).catch((err) => { logger('error', '[MediaEngine] recvTransport connect 失败', err) errback(err) }) }) transport.on('connectionstatechange', (state) => { logger('debug', `[MediaEngine] recvTransport state=${state}`) }) } /** in-flight Promise:sendTransport 创建并发锁;避免突发并发下重复创建 */ let sendTransportPromise = null /** in-flight Promise:recvTransport 创建并发锁;Task 15 后入者 burst 订阅时尤其关键 */ let recvTransportPromise = null /** 按需创建 sendTransport(已存在则复用,并发调用共享同一个 in-flight Promise) */ const ensureSendTransport = async () => { ensureDeviceLoaded() if (sendTransport && !sendTransport.closed) { return sendTransport } if (sendTransportPromise) { return sendTransportPromise } // Nit(代码审查 2026-04-23 第 15 条):in-flight 锁必须在 resolve/reject 两种结局下均置空, // 否则失败后的下次调用会复用 rejected Promise 直接 throw。这里使用 finally 同时覆盖两路。 sendTransportPromise = (async () => { try { return await _createSendTransport() } finally { sendTransportPromise = null } })() return sendTransportPromise } /** 首次创建 sendTransport 的底层流程,原本 inline 在 ensureSendTransport 中 */ const _createSendTransport = async () => { const info = await requestTransportInfo('send') // const iceServers = [ // { // urls: 'turn:123.6.102.114:3478?transport=udp', // username: 'echochat', // credential: 'echochat_2026-1s32dswW@#' // } //]; sendTransport = device.createSendTransport({ id: info.id, iceParameters: info.iceParameters, iceCandidates: info.iceCandidates, dtlsParameters: info.dtlsParameters, sctpParameters: info.sctpParameters, iceServers }) bindSendTransportEvents(sendTransport) logger('info', '[MediaEngine] sendTransport 创建成功', { id: sendTransport.id }) console.log('[MediaEngine] iceParameters:', sendTransport.iceParameters); console.log('[MediaEngine] iceCandidates:', sendTransport.iceCandidates); console.log('[MediaEngine] dtlsParameters:', sendTransport.dtlsParameters); console.log('[MediaEngine] sctpParameters:', sendTransport.sctpParameters); return sendTransport } /** 按需创建 recvTransport(已存在则复用,并发调用共享同一个 in-flight Promise) */ const ensureRecvTransport = async () => { ensureDeviceLoaded() if (recvTransport && !recvTransport.closed) { return recvTransport } if (recvTransportPromise) { return recvTransportPromise } // Nit(同 ensureSendTransport):finally 覆盖 resolve/reject 两路置空 recvTransportPromise = (async () => { try { return await _createRecvTransport() } finally { recvTransportPromise = null } })() return recvTransportPromise } /** 首次创建 recvTransport 的底层流程 */ const _createRecvTransport = async () => { const info = await requestTransportInfo('recv') recvTransport = device.createRecvTransport({ id: info.id, iceParameters: info.iceParameters, iceCandidates: info.iceCandidates, dtlsParameters: info.dtlsParameters, sctpParameters: info.sctpParameters }) bindRecvTransportEvents(recvTransport) logger('info', '[MediaEngine] recvTransport 创建成功', { id: recvTransport.id }) console.log('[MediaEngine] iceParameters:', recvTransport.iceParameters); console.log('[MediaEngine] iceCandidates:', recvTransport.iceCandidates); console.log('[MediaEngine] dtlsParameters:', recvTransport.dtlsParameters); console.log('[MediaEngine] sctpParameters:', recvTransport.sctpParameters); return recvTransport } /** * 在 sendTransport 上创建 Producer(推本地音/视频) * @param {Object} opts * @param {'audio'|'video'} opts.kind * @param {MediaStreamTrack} opts.track * @param {Object} [opts.encodings] * @param {Object} [opts.codecOptions] * @param {Object} [opts.appData] * @returns {Promise} */ const produce = async ({ kind, track, encodings, codecOptions, appData }) => { ensureDeviceLoaded() if (!device.canProduce(kind)) { throw new Error(`[MediaEngine] Device 不支持 produce kind=${kind}`) } const transport = await ensureSendTransport() const produceOpts = { track } if (encodings) produceOpts.encodings = encodings if (codecOptions) produceOpts.codecOptions = codecOptions produceOpts.appData = { user_id: userId, ...(appData || {}) } const producer = await transport.produce(produceOpts) producers.set(producer.id, producer) // 监听 Producer 关闭事件,清理本地索引;外部可见的关闭由 closeProducer 触发 WS producer.on('transportclose', () => { producers.delete(producer.id) logger('debug', `[MediaEngine] producer ${producer.id} transportclose`) }) producer.on('trackended', () => { logger('warn', `[MediaEngine] producer ${producer.id} trackended,将触发关闭`) closeProducer(producer.id).catch((err) => logger('error', '关闭 Producer 失败', err)) }) logger('info', `[MediaEngine] producer 创建成功 kind=${kind} id=${producer.id}`) return producer } /** * 订阅远端 Producer:请求后端创建 Consumer → 本地 consume → 等 track 挂好后 resume * * 与 Task 9 决策 Q6=B 对齐的规范流程: * 1. WS consume.start → 拿到 { id, producerId, kind, rtpParameters } (Node 侧 paused) * 2. recvTransport.consume(...) 得到本地 Consumer(track 可用) * 3. 调用方把 track 挂到