WebRTC 解决的是实时媒体传输,却不会自动处理用户插话、异步响应竞争和缓冲区失控。本文用一个可运行的 TypeScript 状态机,把 VAD、转写、模型、TTS 与播放串成可取消、可背压、可恢复的完整流程。
WebRTC 之后,真正麻烦的是并发
一个语音智能体通常同时运行五条链路:麦克风采集、语音活动检测(VAD)、增量转写、模型生成、语音合成与播放。它们的速度并不一致,也不会严格按顺序结束。
最典型的场景是:智能体正在播放第 7 轮回答,用户突然开口。此时至少要同时完成四件事:停止播放、取消模型与 TTS、清空尚未播放的音频,并确保第 7 轮迟到的数据不能混入第 8 轮。
只设置一个 isSpeaking 布尔值不够,因为它无法表达响应属于哪个轮次,也无法阻止已经在网络或任务队列中的回调继续写入缓冲区。
| 问题 | 仅靠 WebRTC 的结果 | 状态机的处理方式 |
|---|---|---|
| 用户插话 | 旧音频可能继续播放 | 中止任务并立即清空播放队列 |
| 旧响应晚到 | 混入新一轮回答 | 用单调递增的 turnId 丢弃 |
| TTS 快于播放 | 内存与延迟持续增长 | 高水位限制,生产者等待 |
| 模型或播放异常 | 多个组件状态不一致 | 进入恢复态,统一清理资源 |
这里把传输层和控制层分开:WebRTC 负责送入音频;状态机只消费“检测到用户开口”“最终转写完成”等语义事件。这样即使以后换成 WebSocket 或本地音频设备,轮次控制也不需要重写。
定义状态、轮次与取消边界
示例采用四个状态:
listening:等待用户输入,也作为一轮结束后的稳定状态。thinking:模型已经启动,尚无可播放音频。speaking:TTS 音频正在生成、排队或播放。recovering:发生异常,正在取消任务并清理资源。
轮次标识 turnId 必须单调递增。用户一开口就增加轮次,而不是等转写完成后再增加,因为插话发生时首先要让旧响应失效。每个异步边界前后都应检查轮次;单独调用 AbortController.abort() 仍不够,因为某些远端请求可能已经返回,或者第三方实现并不及时响应取消信号。
背压则使用按字节计算的高水位。队列超过上限时,TTS 生产流程等待播放端释放容量,而不是无限追加数据。生产者等待期间仍要监听取消信号,否则插话会被背压等待阻塞。
可运行的 TypeScript 实现
下面的程序不依赖浏览器,可直接观察插话、过期响应丢弃和队列背压。准备 Node.js 20 或更新版本后运行:
npm init -y
npm install -D typescript tsx @types/node
npx tsx index.ts
将以下内容保存为 index.ts:
type Phase = 'listening' | 'thinking' | 'speaking' | 'recovering';
type AudioChunk = {
pcm: Uint8Array;
durationMs: number;
};
type QueuedChunk = AudioChunk & { turn: number };
type Ports = {
model(text: string, signal: AbortSignal): AsyncIterable<string>;
synthesize(text: string, signal: AbortSignal): Promise<AudioChunk>;
play(chunk: AudioChunk, signal: AbortSignal): Promise<void>;
stopPlayback(): void;
};
class FlowAbortError extends Error {
name = 'AbortError';
}
function sleep(ms: number, signal?: AbortSignal): Promise<void> {
return new Promise((resolve, reject) => {
if (signal?.aborted) return reject(new FlowAbortError());
const timer = setTimeout(done, ms);
const onAbort = () => {
clearTimeout(timer);
cleanup();
reject(new FlowAbortError());
};
const cleanup = () => signal?.removeEventListener('abort', onAbort);
function done() {
cleanup();
resolve();
}
signal?.addEventListener('abort', onAbort, { once: true });
});
}
function isAbort(error: unknown): boolean {
return error instanceof Error && error.name === 'AbortError';
}
class VoiceAgent {
private phase: Phase = 'listening';
private turnId = 0;
private generationAbort = new AbortController();
private playbackAbort = new AbortController();
private queue: QueuedChunk[] = [];
private queuedBytes = 0;
private generationDone = false;
private playingTurn?: number;
private closed = false;
constructor(
private readonly ports: Ports,
private readonly highWaterMark = 24_000
) {
void this.playbackLoop();
}
userSpeechStart(): void {
this.turnId += 1;
this.generationAbort.abort();
this.playbackAbort.abort();
this.ports.stopPlayback();
this.clearQueue();
this.generationDone = false;
this.move('listening');
console.log(`[turn ${this.turnId}] 用户开口,旧响应已取消`);
}
finalTranscript(text: string): void {
if (!text.trim() || this.phase !== 'listening') return;
const turn = this.turnId;
this.generationAbort = new AbortController();
this.move('thinking');
void this.generate(turn, text, this.generationAbort.signal);
}
close(): void {
this.closed = true;
this.generationAbort.abort();
this.playbackAbort.abort();
this.ports.stopPlayback();
this.clearQueue();
}
private async generate(
turn: number,
text: string,
signal: AbortSignal
): Promise<void> {
try {
for await (const token of this.ports.model(text, signal)) {
this.assertCurrent(turn, signal);
const audio = await this.ports.synthesize(token, signal);
this.assertCurrent(turn, signal);
this.move('speaking');
await this.enqueue({ ...audio, turn }, signal);
}
this.generationDone = true;
this.finishIfDrained(turn);
} catch (error) {
if (!isAbort(error)) await this.recover(turn, error);
}
}
private async enqueue(item: QueuedChunk, signal: AbortSignal): Promise<void> {
if (item.pcm.byteLength > this.highWaterMark) {
throw new Error('单个音频块超过队列高水位');
}
let reported = false;
while (this.queuedBytes + item.pcm.byteLength > this.highWaterMark) {
if (!reported) console.log('[backpressure] TTS 等待播放端释放容量');
reported = true;
await sleep(10, signal);
this.assertCurrent(item.turn, signal);
}
this.queue.push(item);
this.queuedBytes += item.pcm.byteLength;
}
private async playbackLoop(): Promise<void> {
while (!this.closed) {
const item = this.queue.shift();
if (!item) {
await sleep(5);
continue;
}
this.queuedBytes -= item.pcm.byteLength;
if (item.turn !== this.turnId) continue;
this.playingTurn = item.turn;
this.playbackAbort = new AbortController();
try {
await this.ports.play(item, this.playbackAbort.signal);
} catch (error) {
if (!isAbort(error)) await this.recover(item.turn, error);
} finally {
this.playingTurn = undefined;
}
this.finishIfDrained(item.turn);
}
}
private finishIfDrained(turn: number): void {
if (
turn === this.turnId &&
this.generationDone &&
this.queue.length === 0 &&
this.playingTurn !== turn &&
this.phase !== 'recovering'
) {
this.move('listening');
}
}
private async recover(turn: number, error: unknown): Promise<void> {
if (turn !== this.turnId) return;
console.error(`[turn ${turn}] 异常,开始恢复`, error);
this.move('recovering');
this.generationAbort.abort();
this.playbackAbort.abort();
this.ports.stopPlayback();
this.clearQueue();
await sleep(100);
if (turn === this.turnId && !this.closed) this.move('listening');
}
private assertCurrent(turn: number, signal: AbortSignal): void {
if (signal.aborted || turn !== this.turnId) throw new FlowAbortError();
}
private clearQueue(): void {
this.queue = [];
this.queuedBytes = 0;
}
private move(next: Phase): void {
if (this.phase === next) return;
console.log(`[state] ${this.phase} -> ${next}`);
this.phase = next;
}
}
const ports: Ports = {
async *model(text, signal) {
for (const token of Array.from(`收到:${text}。这是分段生成的回答。`)) {
await sleep(45, signal);
yield token;
}
},
async synthesize(text, signal) {
await sleep(25, signal);
const pcm = new Uint8Array(12_000);
pcm.fill(text.charCodeAt(0) % 255);
return { pcm, durationMs: 160 };
},
async play(chunk, signal) {
await sleep(chunk.durationMs, signal);
console.log(`[play] ${chunk.pcm.byteLength} bytes`);
},
stopPlayback() {
console.log('[play] stop');
}
};
async function main() {
const agent = new VoiceAgent(ports);
agent.userSpeechStart();
agent.finalTranscript('请介绍状态机');
await sleep(260);
agent.userSpeechStart();
agent.finalTranscript('先回答我刚才的插话');
await sleep(2_500);
agent.close();
}
void main();
运行时,第一轮尚未播放完便会被第二次 userSpeechStart 取消。由于每个音频块都携带轮次,即使旧模型稍后返回,也无法进入当前队列。示例故意让合成速度快于播放速度,因此还能看到背压日志。
接入真实音频链路与异常恢复
浏览器端可以用 getUserMedia 获得麦克风轨道,再通过 RTCPeerConnection.addTrack 传输;如果需要本地 PCM 分析,可使用 AudioWorklet。这些 API 负责媒体,不应直接修改模型或播放状态,而应转换为状态机事件:
// VAD 检测到语音起点
agent.userSpeechStart();
// ASR 确认一句最终文本;增量文本只用于界面展示
agent.finalTranscript(finalText);
生产环境要替换示例中的三个端口:model 返回可取消的增量文本,synthesize 返回固定时长的小音频块,play 把 PCM 交给 Web Audio 或原生播放器。stopPlayback 必须真正停止当前音源;仅清空尚未播放的数组,不能停止已经提交给音频设备的数据。
音频块不宜无限大。较小的块更容易快速打断,但会增加调度开销;较大的块调用次数少,插话延迟却可能更明显。合适大小取决于编码格式、播放设备和网络条件,应通过端到端观测调整,而不是只看模型首字延迟。
异常处理也应遵守同一条原则:先确认异常仍属于当前轮次,再进入 recovering。恢复过程统一取消生成、停止播放、清空缓冲,最后回到稳定的监听态。对于网络重连,可以在恢复态之外增加有限次数和退避间隔;不要在多个回调中分别重试,否则容易重复创建模型请求。
总结
实时语音智能体的核心并不是把 WebRTC、ASR、模型和 TTS 依次接通,而是控制它们之间不可避免的并发:
- 用户开口时立即递增
turnId,使旧轮次整体失效。 AbortController用于主动取消,轮次检查用于拦截迟到结果,两者不能互相替代。- 播放队列设置按字节计算的高水位,让快速生产端等待慢速消费端。
- 打断时既要清队列,也要停止正在播放的音源。
- 所有异常通过恢复态统一收敛,避免模型、TTS 与播放器各自重试。
先把这些控制规则做成独立状态机,再接入 WebRTC 和具体云服务,通常更容易测试,也更容易定位插话失效、串轮和延迟不断累积的问题。
评论