inspection-host-computer/server/audioWebSocket.js
2026-07-22 13:33:59 +08:00

171 lines
5.9 KiB
JavaScript
Raw Permalink 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.

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
}