From b54e0049c5f4bb2850bdbcef5e79801eb90036fc Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Fri, 31 Jul 2026 15:56:09 +0800 Subject: [PATCH] fix(ui): bound realtime Talk PCM playback ownership --- ui/src/pages/chat/realtime-talk-audio.test.ts | 149 +++++++++++++++++- ui/src/pages/chat/realtime-talk-audio.ts | 50 +++++- 2 files changed, 190 insertions(+), 9 deletions(-) diff --git a/ui/src/pages/chat/realtime-talk-audio.test.ts b/ui/src/pages/chat/realtime-talk-audio.test.ts index a1e1ccad9e07..b0db4684e9b6 100644 --- a/ui/src/pages/chat/realtime-talk-audio.test.ts +++ b/ui/src/pages/chat/realtime-talk-audio.test.ts @@ -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(); + }); +}); diff --git a/ui/src/pages/chat/realtime-talk-audio.ts b/ui/src/pages/chat/realtime-talk-audio.ts index bdb7aef45d81..32590083661b 100644 --- a/ui/src/pages/chat/realtime-talk-audio.ts +++ b/ui/src/pages/chat/realtime-talk-audio.ts @@ -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(); @@ -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; } }