inspection-host-computer/server/audioWebSocket.js

171 lines
5.9 KiB
JavaScript
Raw Permalink Normal View History

2026-07-22 13:33:59 +08:00
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
}