From a84ede72489ca01e67759ec49a90b5d8837a18cd Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Fri, 31 Jul 2026 13:27:06 +0800 Subject: [PATCH] fix(ui): bound Talk relay microphone uplink --- .../chat/realtime-talk-gateway-relay.test.ts | 110 ++++++++++++++++++ .../pages/chat/realtime-talk-gateway-relay.ts | 50 ++++++-- 2 files changed, 152 insertions(+), 8 deletions(-) diff --git a/ui/src/pages/chat/realtime-talk-gateway-relay.test.ts b/ui/src/pages/chat/realtime-talk-gateway-relay.test.ts index 35e970bf7be0..c9c206720cac 100644 --- a/ui/src/pages/chat/realtime-talk-gateway-relay.test.ts +++ b/ui/src/pages/chat/realtime-talk-gateway-relay.test.ts @@ -380,6 +380,116 @@ describe("GatewayRelayRealtimeTalkTransport", () => { expect(onInputLevel).toHaveBeenLastCalledWith(0); }); + it("bounds stalled microphone appends and aborts every owner on stop", async () => { + const onStatus = vi.fn(); + const client = createClient(); + let activeAppends = 0; + let peakActiveAppends = 0; + const appendSignals: AbortSignal[] = []; + vi.mocked(client["request"]).mockImplementation((method, _params, options) => { + if (method !== "talk.session.appendAudio") { + return Promise.resolve({}); + } + const signal = options?.signal; + if (!signal) { + return Promise.reject(new Error("missing append abort signal")); + } + appendSignals.push(signal); + activeAppends += 1; + peakActiveAppends = Math.max(peakActiveAppends, activeAppends); + return new Promise((_, reject) => { + signal.addEventListener( + "abort", + () => { + activeAppends -= 1; + reject(new Error("append aborted")); + }, + { once: true }, + ); + }); + }); + const transport = createTransport({ callbacks: { onStatus }, client }); + + await transport.start(); + const samples = new Float32Array(4096); + for (let index = 0; index < 10_000; index += 1) { + pumpMicrophone(samples); + } + + const appendCalls = requestCallsFor(client, "talk.session.appendAudio"); + expect(appendCalls).toHaveLength(4); + expect(peakActiveAppends).toBe(4); + expect(activeAppends).toBe(4); + expect(new Set(appendSignals).size).toBe(1); + expect( + appendCalls.every( + (call) => call[2]?.signal === appendSignals[0] && call[2]?.timeoutMs === 8_000, + ), + ).toBe(true); + + transport.stop(); + transport.stop(); + await Promise.resolve(); + + expect(activeAppends).toBe(0); + expect(appendSignals.every((signal) => signal.aborted)).toBe(true); + expect(requestCallsFor(client, "talk.session.close")).toHaveLength(1); + expect(onStatus).not.toHaveBeenCalled(); + }); + + it("preserves accepted microphone frame order", async () => { + const client = createClient(); + const transport = createTransport({ client }); + + await transport.start(); + for (const timestamp of [10, 20, 30, 40]) { + audioCurrentTime = timestamp / 1_000; + pumpMicrophone(new Float32Array(4096)); + } + + expect( + requestCallsFor(client, "talk.session.appendAudio").map( + (call) => (call[1] as { timestamp: number }).timestamp, + ), + ).toEqual([10, 20, 30, 40]); + transport.stop(); + }); + + it("ignores a stale append rejection after a replacement starts", async () => { + const oldStatus = vi.fn(); + const oldClient = createClient(); + let rejectOldAppend: (error: Error) => void = () => undefined; + vi.mocked(oldClient["request"]).mockImplementation((method) => { + if (method !== "talk.session.appendAudio") { + return Promise.resolve({}); + } + return new Promise((_, reject) => { + rejectOldAppend = reject; + }); + }); + const oldTransport = createTransport({ callbacks: { onStatus: oldStatus }, client: oldClient }); + + await oldTransport.start(); + pumpMicrophone(new Float32Array(4096)); + oldTransport.stop(); + + const replacementStatus = vi.fn(); + const replacementClient = createClient(); + const replacement = createTransport({ + callbacks: { onStatus: replacementStatus }, + client: replacementClient, + }); + await replacement.start(); + pumpMicrophone(new Float32Array(4096)); + rejectOldAppend(new Error("late stale append failure")); + await Promise.resolve(); + + expect(requestCallsFor(replacementClient, "talk.session.appendAudio")).toHaveLength(1); + expect(oldStatus).not.toHaveBeenCalled(); + expect(replacementStatus).not.toHaveBeenCalled(); + replacement.stop(); + }); + it("stops microphone pumping when the relay rejects appended audio", async () => { const onStatus = vi.fn(); const client = createClient(); diff --git a/ui/src/pages/chat/realtime-talk-gateway-relay.ts b/ui/src/pages/chat/realtime-talk-gateway-relay.ts index 747f07b4cd35..b71189dfedf1 100644 --- a/ui/src/pages/chat/realtime-talk-gateway-relay.ts +++ b/ui/src/pages/chat/realtime-talk-gateway-relay.ts @@ -22,6 +22,8 @@ import { const BARGE_IN_RMS_THRESHOLD = 0.02; const BARGE_IN_PEAK_THRESHOLD = 0.08; const BARGE_IN_CONSECUTIVE_SPEECH_FRAMES = 2; +const MAX_PENDING_AUDIO_APPENDS = 4; +const AUDIO_APPEND_TIMEOUT_MS = 8_000; export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport { private media: MediaStream | null = null; @@ -31,6 +33,8 @@ export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport private readonly inputPump = new RealtimeTalkPcmInputPump(); private unsubscribe: (() => void) | null = null; private closed = false; + private audioAppendAbortController: AbortController | null = null; + private readonly pendingAudioAppends = new Set>(); private readonly outputQueue = new RealtimeTalkPcmOutputQueue(); private readonly consultAbortControllers = new Map(); private readonly completedToolCalls = new Set(); @@ -80,6 +84,8 @@ export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport this.media = media; this.inputContext = new AudioContext({ sampleRate: this.session.audio.inputSampleRateHz }); this.outputContext = new AudioContext({ sampleRate: this.session.audio.outputSampleRateHz }); + this.abortPendingAudioAppends(); + this.audioAppendAbortController = new AbortController(); if (this.ctx.callbacks.onInputLevel) { this.inputMeter = new RealtimeTalkMediaStreamMeter(this.ctx.callbacks.onInputLevel); this.inputMeter.start(this.media, this.inputContext); @@ -104,6 +110,7 @@ export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport this.unsubscribe?.(); this.unsubscribe = null; this.inputPump.stop(); + this.abortPendingAudioAppends(); this.inputMeter?.stop(); this.inputMeter = null; // Mark callbacks recurse until playback drains, so shutdown must cancel every owned timer. @@ -128,18 +135,35 @@ export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport if (this.closed) { return; } - const pcm = floatToPcm16(samples); if (this.detectBargeInSpeech(samples)) { this.cancelOutputForBargeIn(); } - void this.ctx.client - .request("talk.session.appendAudio", { - sessionId: this.session.relaySessionId, - audioBase64: bytesToBase64(pcm), - timestamp: Math.round((this.inputContext?.currentTime ?? 0) * 1000), - }) + const abortController = this.audioAppendAbortController; + // Live microphone frames become stale once the Gateway falls behind, so drop new + // frames at the ownership cap instead of growing a latency queue. + if ( + !abortController || + abortController.signal.aborted || + this.pendingAudioAppends.size >= MAX_PENDING_AUDIO_APPENDS + ) { + return; + } + const pcm = floatToPcm16(samples); + const request = this.ctx.client + .request( + "talk.session.appendAudio", + { + sessionId: this.session.relaySessionId, + audioBase64: bytesToBase64(pcm), + timestamp: Math.round((this.inputContext?.currentTime ?? 0) * 1000), + }, + { + signal: abortController.signal, + timeoutMs: AUDIO_APPEND_TIMEOUT_MS, + }, + ) .catch((error: unknown) => { - if (!this.closed) { + if (!this.closed && !abortController.signal.aborted) { this.ctx.callbacks.onStatus?.( "error", error instanceof Error ? error.message : String(error), @@ -147,9 +171,19 @@ export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport this.stop(); } }); + this.pendingAudioAppends.add(request); + void request.finally(() => { + this.pendingAudioAppends.delete(request); + }); }); } + private abortPendingAudioAppends(): void { + this.audioAppendAbortController?.abort(); + this.audioAppendAbortController = null; + this.pendingAudioAppends.clear(); + } + private handleRelayEvent(event: GatewayRelayEvent): void { if (event.relaySessionId !== this.session.relaySessionId || this.closed) { return;