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 和具体云服务,通常更容易测试,也更容易定位插话失效、串轮和延迟不断累积的问题。