fix(ui): bound realtime Talk PCM playback ownership

This commit is contained in:
Vincent Koc
2026-07-31 15:56:09 +08:00
parent a5c8c5b3ec
commit b54e0049c5
2 changed files with 190 additions and 9 deletions
+148 -1
View File
@@ -1,6 +1,52 @@
// @vitest-environment jsdom
import { afterEach, describe, expect, it, vi } from "vitest";
import { RealtimeTalkMediaStreamMeter } from "./realtime-talk-audio.ts";
import {
bytesToBase64,
RealtimeTalkMediaStreamMeter,
RealtimeTalkPcmOutputQueue,
} from "./realtime-talk-audio.ts";
class MockAudioBufferSource {
buffer: unknown = null;
readonly connect = vi.fn();
readonly start = vi.fn();
readonly stop = vi.fn();
private ended: (() => void) | null = null;
addEventListener(type: string, handler: () => void): void {
if (type === "ended") {
this.ended = handler;
}
}
emitEnded(): void {
this.ended?.();
}
}
class MockOutputAudioContext {
currentTime = 0;
readonly destination = {};
readonly sources: MockAudioBufferSource[] = [];
createBuffer(_channels: number, length: number, sampleRate: number) {
const channel = new Float32Array(length);
return {
duration: length / sampleRate,
getChannelData: () => channel,
};
}
createBufferSource(): MockAudioBufferSource {
const source = new MockAudioBufferSource();
this.sources.push(source);
return source;
}
}
function silentPcmBase64(sampleCount: number): string {
return bytesToBase64(new Uint8Array(sampleCount * 2));
}
describe("RealtimeTalkMediaStreamMeter", () => {
afterEach(() => {
@@ -66,3 +112,104 @@ describe("RealtimeTalkMediaStreamMeter", () => {
expect(onLevel).toHaveBeenLastCalledWith(0);
});
});
describe("RealtimeTalkPcmOutputQueue", () => {
it("preserves ordered playback while the AudioContext advances normally", () => {
const context = new MockOutputAudioContext();
context.currentTime = 1;
const queue = new RealtimeTalkPcmOutputQueue();
expect(queue.play(silentPcmBase64(100), context as unknown as AudioContext, 100)).toBe(
"queued",
);
context.currentTime = 1.5;
expect(queue.play(silentPcmBase64(50), context as unknown as AudioContext, 100)).toBe("queued");
expect(context.sources.map((source) => source.start.mock.calls[0]?.[0])).toEqual([1, 2]);
expect(queue.queuedUntil).toBe(2.5);
expect(queue.isPlaying).toBe(true);
});
it("bounds a frozen AudioContext by queued seconds before allocating another source", () => {
const context = new MockOutputAudioContext();
const queue = new RealtimeTalkPcmOutputQueue();
expect(queue.play(silentPcmBase64(600), context as unknown as AudioContext, 100)).toBe(
"queued",
);
expect(queue.play(silentPcmBase64(500), context as unknown as AudioContext, 100)).toBe(
"overflow",
);
expect(context.sources).toHaveLength(1);
expect(queue.queuedUntil).toBe(6);
});
it("rejects an oversized frame before base64 decoding", () => {
const context = new MockOutputAudioContext();
const queue = new RealtimeTalkPcmOutputQueue();
expect(queue.play("!".repeat(3_000), context as unknown as AudioContext, 100)).toBe("overflow");
expect(context.sources).toHaveLength(0);
});
it("hard-caps source ownership across ten thousand suspended-context chunks", () => {
const context = new MockOutputAudioContext();
const queue = new RealtimeTalkPcmOutputQueue();
let queued = 0;
let overflowed = 0;
for (let index = 0; index < 10_000; index += 1) {
const result = queue.play(silentPcmBase64(1), context as unknown as AudioContext, 48_000);
if (result === "queued") {
queued += 1;
} else if (result === "overflow") {
overflowed += 1;
}
}
expect(queued).toBe(320);
expect(overflowed).toBe(9_680);
expect(context.sources).toHaveLength(320);
});
it("releases source ownership on ended", () => {
const context = new MockOutputAudioContext();
const queue = new RealtimeTalkPcmOutputQueue();
const chunk = silentPcmBase64(1);
for (let index = 0; index < 320; index += 1) {
expect(queue.play(chunk, context as unknown as AudioContext, 48_000)).toBe("queued");
}
expect(queue.play(chunk, context as unknown as AudioContext, 48_000)).toBe("overflow");
context.sources[0]?.emitEnded();
expect(queue.play(chunk, context as unknown as AudioContext, 48_000)).toBe("queued");
expect(context.sources).toHaveLength(321);
});
it("stops idempotently and isolates late ended events from replacement playback", () => {
const context = new MockOutputAudioContext();
const queue = new RealtimeTalkPcmOutputQueue();
const chunk = silentPcmBase64(100);
expect(queue.play(chunk, context as unknown as AudioContext, 100)).toBe("queued");
const oldSource = context.sources[0];
context.currentTime = 0.25;
queue.stop(context as unknown as AudioContext);
queue.stop(context as unknown as AudioContext);
expect(oldSource?.stop).toHaveBeenCalledOnce();
expect(queue.isPlaying).toBe(false);
expect(queue.queuedUntil).toBe(0.25);
expect(queue.play(chunk, context as unknown as AudioContext, 100)).toBe("queued");
const replacementSource = context.sources[1];
oldSource?.emitEnded();
expect(queue.isPlaying).toBe(true);
expect(queue.queuedUntil).toBe(1.25);
expect(replacementSource?.stop).not.toHaveBeenCalled();
});
});
+42 -8
View File
@@ -214,6 +214,16 @@ function pcm16ToFloat(bytes: Uint8Array): Float32Array {
return samples;
}
function base64DecodedByteLength(value: string): number {
const padding = value.endsWith("==") ? 2 : value.endsWith("=") ? 1 : 0;
return Math.max(0, Math.floor((value.length * 3) / 4) - padding);
}
const REALTIME_TALK_PCM_OUTPUT_MAX_QUEUED_SECONDS = 10;
const REALTIME_TALK_PCM_OUTPUT_MAX_SOURCES = 320;
type RealtimeTalkPcmOutputQueuePlayResult = "queued" | "ignored" | "overflow";
export class RealtimeTalkPcmOutputQueue {
private playhead = 0;
private readonly sources = new Set<AudioBufferSourceNode>();
@@ -226,13 +236,34 @@ export class RealtimeTalkPcmOutputQueue {
return this.sources.size > 0;
}
play(base64: string, outputContext: AudioContext | null, outputSampleRateHz: number): void {
play(
base64: string,
outputContext: AudioContext | null,
outputSampleRateHz: number,
): RealtimeTalkPcmOutputQueuePlayResult {
if (!outputContext) {
return;
return "ignored";
}
const startAt = Math.max(outputContext.currentTime, this.playhead);
const queuedSeconds = Math.max(0, startAt - outputContext.currentTime);
const remainingSeconds = REALTIME_TALK_PCM_OUTPUT_MAX_QUEUED_SECONDS - queuedSeconds;
const decodedByteLength = base64DecodedByteLength(base64);
const sampleCount = Math.floor(decodedByteLength / 2);
if (
this.sources.size >= REALTIME_TALK_PCM_OUTPUT_MAX_SOURCES ||
remainingSeconds <= 0 ||
sampleCount / outputSampleRateHz > remainingSeconds
) {
return "overflow";
}
const samples = pcm16ToFloat(base64ToBytes(base64));
if (samples.length === 0) {
return;
return "ignored";
}
const duration = samples.length / outputSampleRateHz;
const nextPlayhead = startAt + duration;
if (nextPlayhead - outputContext.currentTime > REALTIME_TALK_PCM_OUTPUT_MAX_QUEUED_SECONDS) {
return "overflow";
}
const buffer = outputContext.createBuffer(1, samples.length, outputSampleRateHz);
buffer.getChannelData(0).set(samples);
@@ -241,18 +272,21 @@ export class RealtimeTalkPcmOutputQueue {
source.addEventListener("ended", () => this.sources.delete(source));
source.buffer = buffer;
source.connect(outputContext.destination);
const startAt = Math.max(outputContext.currentTime, this.playhead);
source.start(startAt);
this.playhead = startAt + buffer.duration;
this.playhead = nextPlayhead;
return "queued";
}
stop(outputContext: AudioContext | null): void {
for (const source of this.sources) {
// Release ownership first so synchronous or late `ended` events from stopped
// sources cannot affect audio queued by a replacement playback turn.
const sources = [...this.sources];
this.sources.clear();
this.playhead = outputContext?.currentTime ?? 0;
for (const source of sources) {
try {
source.stop();
} catch {}
}
this.sources.clear();
this.playhead = outputContext?.currentTime ?? 0;
}
}