import { WebSocketServer, WebSocket } from 'ws' import { createMicrophoneGrpcClient, createSpeakerGrpcClient } from './grpcClient.js' // 设备实际返回 48kHz;上下行统一使用 48kHz,避免设备按固定采样率解释 44.1kHz 数据。 const SAMPLE_RATE = 48000 const CHANNELS = 2 const MAX_BUFFERED_BYTES = SAMPLE_RATE * CHANNELS * 2 * 2 function sendJson(socket, payload) { if (socket.readyState === WebSocket.OPEN) socket.send(JSON.stringify(payload)) } function errorText(error) { return error?.details || error?.message || '机器人音频服务异常' } function waitForGrpcReady(client, timeoutMs = 5000) { return new Promise((resolve, reject) => { client.waitForReady(Date.now() + timeoutMs, (error) => { if (error) reject(error) else resolve() }) }) } export function attachAudioWebSocket(httpServer) { const wss = new WebSocketServer({ server: httpServer, path: '/audio' }) wss.on('connection', (socket) => { let initialized = false let speakerId = '' let microphoneStream let speakerStream let microphoneClient let speakerClient let remoteSampleRate = SAMPLE_RATE // 设为 0,确保首帧一定向浏览器发送真实声道数。 let remoteChannels = 0 let hasLoggedRemoteFormat = false let grpcReady = false let hasFailed = false let closing = false const closeGrpc = () => { closing = true microphoneStream?.cancel() microphoneStream = undefined speakerStream?.end() speakerStream = undefined microphoneClient?.close() speakerClient?.close() microphoneClient = undefined speakerClient = undefined } const fail = (error) => { if (hasFailed || closing) return hasFailed = true console.error('[audio] gRPC error:', error) sendJson(socket, { type: grpcReady ? 'error' : 'connection-failed', message: grpcReady ? errorText(error) : `无法连接机器人音频服务:${errorText(error)}` }) closeGrpc() } socket.on('message', (data, isBinary) => { if (isBinary) { if (!initialized || !speakerStream || speakerStream.destroyed) return const pcm = Buffer.from(data) if (!pcm.length || pcm.length % 2 !== 0) return speakerStream.write({ header: { device_id: speakerId }, audio: { data: pcm, sample_rate: SAMPLE_RATE, channels: CHANNELS, format: 'PCM', codec: 'pcm_s16le', // nb_samples 表示每个声道的采样帧数,不是所有声道样本总数。 nb_samples: pcm.length / (2 * CHANNELS) } }) return } let message try { message = JSON.parse(data.toString()) } catch { sendJson(socket, { type: 'error', message: '无效的控制消息' }) return } if (message.type === 'stop') { closeGrpc() socket.close(1000, 'call ended') return } if (message.type !== 'init' || initialized) return const { microphoneId, robotAddress } = message speakerId = message.speakerId?.trim() if (![microphoneId, speakerId, robotAddress].every((value) => typeof value === 'string' && value.trim())) { sendJson(socket, { type: 'error', message: '机器人地址、麦克风 ID 和扬声器 ID 不能为空' }) socket.close(1008, 'invalid init') return } initialized = true microphoneClient = createMicrophoneGrpcClient(robotAddress.trim()) speakerClient = createSpeakerGrpcClient(robotAddress.trim()) // 只有两条 gRPC 通道都真正可用后,前端才进入“通话中”状态。 Promise.all([ waitForGrpcReady(microphoneClient), waitForGrpcReady(speakerClient) ]).then(() => { if (hasFailed || socket.readyState !== WebSocket.OPEN) return grpcReady = true sendJson(socket, { type: 'ready', sampleRate: SAMPLE_RATE, channels: CHANNELS }) }).catch(fail) // 现有 proto 用一条服务端流和一条客户端流共同组成全双工音频通道。 microphoneStream = microphoneClient.streamAudio({ header: { device_id: microphoneId.trim() } }) speakerStream = speakerClient.streamAudio((error, response) => { if (error) fail(error) else if (response?.header?.success === false) { fail(new Error(response.header.error_message || '机器人扬声器拒绝音频流')) } }) microphoneStream.on('data', (frame) => { const audio = frame?.audio if (!audio?.data?.length || socket.readyState !== WebSocket.OPEN) return if (!hasLoggedRemoteFormat) { hasLoggedRemoteFormat = true } const frameSampleRate = Number(audio.sample_rate) || SAMPLE_RATE const frameChannels = Number(audio.channels) || 1 if (frameSampleRate !== remoteSampleRate || frameChannels !== remoteChannels) { remoteSampleRate = frameSampleRate remoteChannels = frameChannels // WebSocket 保证消息顺序,浏览器会先收到格式通知,再收到对应二进制帧。 sendJson(socket, { type: 'audio-format', sampleRate: remoteSampleRate, channels: remoteChannels }) } // 浏览器播放跟不上时丢弃新帧,防止内存持续增长。 if (socket.bufferedAmount <= MAX_BUFFERED_BYTES) socket.send(audio.data, { binary: true }) }) microphoneStream.on('error', (error) => { if (error.code !== 1) fail(error) // CANCELLED(1) 是正常清理 }) microphoneStream.on('end', () => sendJson(socket, { type: 'remote-ended' })) speakerStream.on('error', fail) }) socket.on('close', closeGrpc) socket.on('error', (error) => console.error('[audio] WebSocket error:', error)) }) return wss }