From 0c44438761ef0673bfe22404de82088698d18e5c Mon Sep 17 00:00:00 2001 From: Alix-007 Date: Thu, 16 Jul 2026 21:06:04 +0800 Subject: [PATCH 1/9] fix(ui): stop stalled Google Live connections --- .../chat/realtime-talk-google-live.test.ts | 133 ++++++++++++++++++ .../pages/chat/realtime-talk-google-live.ts | 61 ++++++-- 2 files changed, 181 insertions(+), 13 deletions(-) diff --git a/ui/src/pages/chat/realtime-talk-google-live.test.ts b/ui/src/pages/chat/realtime-talk-google-live.test.ts index 044a0339438a..4cfe54f5d470 100644 --- a/ui/src/pages/chat/realtime-talk-google-live.test.ts +++ b/ui/src/pages/chat/realtime-talk-google-live.test.ts @@ -22,6 +22,7 @@ type MockWebSocketEventType = "close" | "error" | "message" | "open"; const wsInstances: MockGoogleLiveWebSocket[] = []; const createdSources: MockAudioBufferSource[] = []; +const audioContexts: MockAudioContext[] = []; const inputProcessors: Array<{ connect: ReturnType; disconnect: ReturnType; @@ -74,6 +75,19 @@ class MockGoogleLiveWebSocket { } } + emitClose() { + this.readyState = 3; + for (const handler of this.handlers.close) { + handler(); + } + } + + emitError() { + for (const handler of this.handlers.error) { + handler(); + } + } + emitMessage(data: unknown) { for (const handler of this.handlers.message) { handler({ data }); @@ -97,6 +111,7 @@ class MockAudioContext { constructor(options?: { sampleRate?: number }) { this.sampleRate = options?.sampleRate ?? 24000; + audioContexts.push(this); } createMediaStreamSource() { @@ -231,6 +246,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { beforeEach(() => { wsInstances.length = 0; createdSources.length = 0; + audioContexts.length = 0; inputProcessors.length = 0; inputSinks.length = 0; vi.stubGlobal("WebSocket", MockGoogleLiveWebSocket); @@ -246,6 +262,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { }); afterEach(() => { + vi.useRealTimers(); vi.unstubAllGlobals(); }); @@ -336,6 +353,122 @@ describe("GoogleLiveRealtimeTalkTransport", () => { expect(onInputLevel).not.toHaveBeenCalled(); }); + it("times out a stalled WebSocket opening and releases browser resources", async () => { + vi.useFakeTimers(); + const stopTrack = vi.fn(); + getUserMedia.mockResolvedValue({ + getTracks: () => [{ stop: stopTrack }], + }); + const onStatus = vi.fn(); + const onTalkEvent = vi.fn(); + const transport = createTransport({ onStatus, onTalkEvent }); + + await transport.start(); + const ws = latestWebSocket(); + ws.readyState = 0; + await vi.advanceTimersByTimeAsync(30_000); + + expect(onStatus).toHaveBeenCalledExactlyOnceWith( + "error", + "Realtime connection timed out after 30000ms", + ); + expect(stopTrack).toHaveBeenCalledOnce(); + expect(audioContexts).toHaveLength(2); + expect(audioContexts.every((context) => context.close.mock.calls.length === 1)).toBe(true); + expect(ws.readyState).toBe(3); + expect(onTalkEvent).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ type: "session.closed", final: true }), + ); + + await vi.advanceTimersByTimeAsync(30_000); + expect(onStatus).toHaveBeenCalledTimes(1); + }); + + it("clears the WebSocket opening timeout after open", async () => { + vi.useFakeTimers(); + const onStatus = vi.fn(); + const transport = createTransport({ onStatus }); + + await transport.start(); + latestWebSocket().emitOpen(); + const statusCalls = onStatus.mock.calls.slice(); + await vi.advanceTimersByTimeAsync(30_000); + + expect(onStatus.mock.calls).toStrictEqual(statusCalls); + transport.stop(); + }); + + it("releases browser resources once when WebSocket error is followed by close", async () => { + vi.useFakeTimers(); + const stopTrack = vi.fn(); + getUserMedia.mockResolvedValue({ + getTracks: () => [{ stop: stopTrack }], + }); + const onInputLevel = vi.fn(); + const onStatus = vi.fn(); + const onTalkEvent = vi.fn(); + const transport = createTransport({ onInputLevel, onStatus, onTalkEvent }); + + await transport.start(); + const ws = latestWebSocket(); + ws.readyState = 0; + ws.emitError(); + ws.emitClose(); + + expect(onStatus).toHaveBeenCalledExactlyOnceWith("error", "Realtime connection failed"); + expect(stopTrack).toHaveBeenCalledOnce(); + expect(audioContexts).toHaveLength(2); + for (const context of audioContexts) { + expect(context.close).toHaveBeenCalledOnce(); + } + expect(onInputLevel).toHaveBeenLastCalledWith(0); + expect(onTalkEvent).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ type: "session.closed", final: true }), + ); + expect(vi.getTimerCount()).toBe(0); + }); + + it("releases browser resources once when WebSocket closes before open", async () => { + vi.useFakeTimers(); + const stopTrack = vi.fn(); + getUserMedia.mockResolvedValue({ + getTracks: () => [{ stop: stopTrack }], + }); + const onInputLevel = vi.fn(); + const onStatus = vi.fn(); + const onTalkEvent = vi.fn(); + const transport = createTransport({ onInputLevel, onStatus, onTalkEvent }); + + await transport.start(); + const ws = latestWebSocket(); + ws.readyState = 0; + ws.emitClose(); + + expect(onStatus).toHaveBeenCalledExactlyOnceWith("error", "Realtime connection closed"); + expect(stopTrack).toHaveBeenCalledOnce(); + expect(audioContexts).toHaveLength(2); + for (const context of audioContexts) { + expect(context.close).toHaveBeenCalledOnce(); + } + expect(onInputLevel).toHaveBeenLastCalledWith(0); + expect(onTalkEvent).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ type: "session.closed", final: true }), + ); + expect(vi.getTimerCount()).toBe(0); + }); + + it("clears the WebSocket opening timeout when stopped", async () => { + vi.useFakeTimers(); + const onStatus = vi.fn(); + const transport = createTransport({ onStatus }); + + await transport.start(); + transport.stop(); + await vi.advanceTimersByTimeAsync(30_000); + + expect(onStatus).not.toHaveBeenCalled(); + }); + it("requests ArrayBuffer frames and decodes binary setup messages", async () => { const onStatus = vi.fn(); const onTalkEvent = vi.fn(); diff --git a/ui/src/pages/chat/realtime-talk-google-live.ts b/ui/src/pages/chat/realtime-talk-google-live.ts index 12b6ef7bac9b..53c2da8d3dda 100644 --- a/ui/src/pages/chat/realtime-talk-google-live.ts +++ b/ui/src/pages/chat/realtime-talk-google-live.ts @@ -69,6 +69,7 @@ function googleLiveVideoMessage(frame: RealtimeTalkVideoFrame): unknown { }, }; } +const GOOGLE_LIVE_WEBSOCKET_OPEN_TIMEOUT_MS = 30_000; // Browser sessions can still pin a 2.5 model, whose text and tool-response wire // contract differs from the 3.1 default carried in new session metadata. @@ -106,6 +107,7 @@ function buildGoogleLiveUrl(session: RealtimeTalkJsonPcmWebSocketSessionResult): export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { private ws: WebSocket | null = null; + private openTimeout: ReturnType | null = null; private media: MediaStream | null = null; private inputContext: AudioContext | null = null; private outputContext: AudioContext | null = null; @@ -181,27 +183,36 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { this.inputMeter = new RealtimeTalkMediaStreamMeter(this.ctx.callbacks.onInputLevel); this.inputMeter.start(this.media, this.inputContext); } - this.ws = new WebSocket(wsUrl); - this.ws.binaryType = "arraybuffer"; - this.ws.addEventListener("open", () => { - if (this.closed) { + const ws = new WebSocket(wsUrl); + this.ws = ws; + ws.binaryType = "arraybuffer"; + this.openTimeout = globalThis.setTimeout(() => { + if (this.closed || this.ws !== ws) { + return; + } + this.openTimeout = null; + this.ctx.callbacks.onStatus?.( + "error", + `Realtime connection timed out after ${GOOGLE_LIVE_WEBSOCKET_OPEN_TIMEOUT_MS}ms`, + ); + this.stop(); + }, GOOGLE_LIVE_WEBSOCKET_OPEN_TIMEOUT_MS); + ws.addEventListener("open", () => { + this.clearOpenTimeout(ws); + if (this.closed || this.ws !== ws) { return; } this.send(this.session.initialMessage ?? { setup: {} }); this.startMicrophonePump(); }); - this.ws.addEventListener("message", (event) => { + ws.addEventListener("message", (event) => { void this.handleMessage(event.data); }); - this.ws.addEventListener("close", () => { - if (!this.closed) { - this.ctx.callbacks.onStatus?.("error", "Realtime connection closed"); - } + ws.addEventListener("close", () => { + this.failSocket(ws, "Realtime connection closed"); }); - this.ws.addEventListener("error", () => { - if (!this.closed) { - this.ctx.callbacks.onStatus?.("error", "Realtime connection failed"); - } + ws.addEventListener("error", () => { + this.failSocket(ws, "Realtime connection failed"); }); } @@ -221,6 +232,7 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { this.mediaSetupController?.abort(); this.mediaSetupController = null; this.setupComplete = false; + this.clearOpenTimeout(); for (const controller of this.consultAbortControllers) { controller.abort(); } @@ -241,6 +253,29 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { this.ws = null; } + private clearOpenTimeout(ws?: WebSocket): void { + // A stopped socket can emit after a replacement starts. Keep stale events + // from clearing the replacement attempt's startup deadline. + if (ws && this.ws !== ws) { + return; + } + if (this.openTimeout !== null) { + globalThis.clearTimeout(this.openTimeout); + this.openTimeout = null; + } + } + + private failSocket(ws: WebSocket, message: string): void { + // Error and close can arrive for the same socket, or after a replacement + // starts. Only the current lifecycle owner may report and release resources. + if (this.closed || this.ws !== ws) { + return; + } + this.clearOpenTimeout(ws); + this.ctx.callbacks.onStatus?.("error", message); + this.stop(); + } + private startMicrophonePump(): void { if (this.closed || !this.media || !this.inputContext) { return; From efca96c0ece6516a07a6b5da068a4b81695be3c7 Mon Sep 17 00:00:00 2001 From: Alix-007 Date: Fri, 17 Jul 2026 01:25:40 +0800 Subject: [PATCH 2/9] fix(ui): cover Google Live setup with timeout --- .../chat/realtime-talk-google-live.test.ts | 35 ++++++++++++++++--- .../pages/chat/realtime-talk-google-live.ts | 34 +++++++++--------- 2 files changed, 48 insertions(+), 21 deletions(-) diff --git a/ui/src/pages/chat/realtime-talk-google-live.test.ts b/ui/src/pages/chat/realtime-talk-google-live.test.ts index 4cfe54f5d470..57527f8f5057 100644 --- a/ui/src/pages/chat/realtime-talk-google-live.test.ts +++ b/ui/src/pages/chat/realtime-talk-google-live.test.ts @@ -384,17 +384,44 @@ describe("GoogleLiveRealtimeTalkTransport", () => { expect(onStatus).toHaveBeenCalledTimes(1); }); - it("clears the WebSocket opening timeout after open", async () => { + it("times out when WebSocket opens but Google setup never completes", async () => { + vi.useFakeTimers(); + const stopTrack = vi.fn(); + getUserMedia.mockResolvedValue({ + getTracks: () => [{ stop: stopTrack }], + }); + const onStatus = vi.fn(); + const onTalkEvent = vi.fn(); + const transport = createTransport({ onStatus, onTalkEvent }); + + await transport.start(); + latestWebSocket().emitOpen(); + await vi.advanceTimersByTimeAsync(30_000); + + expect(onStatus).toHaveBeenCalledExactlyOnceWith( + "error", + "Realtime connection timed out after 30000ms", + ); + expect(stopTrack).toHaveBeenCalledOnce(); + expect(audioContexts.every((context) => context.close.mock.calls.length === 1)).toBe(true); + expect(onTalkEvent).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ type: "session.closed", final: true }), + ); + }); + + it("clears the startup timeout after Google setup completes", async () => { vi.useFakeTimers(); const onStatus = vi.fn(); const transport = createTransport({ onStatus }); await transport.start(); - latestWebSocket().emitOpen(); - const statusCalls = onStatus.mock.calls.slice(); + const ws = latestWebSocket(); + ws.emitOpen(); + ws.emitMessage(encodeJsonFrame({ setupComplete: {} })); + await flushMicrotasks(); await vi.advanceTimersByTimeAsync(30_000); - expect(onStatus.mock.calls).toStrictEqual(statusCalls); + expect(onStatus).toHaveBeenCalledExactlyOnceWith("listening"); transport.stop(); }); diff --git a/ui/src/pages/chat/realtime-talk-google-live.ts b/ui/src/pages/chat/realtime-talk-google-live.ts index 53c2da8d3dda..428c4f3d10d5 100644 --- a/ui/src/pages/chat/realtime-talk-google-live.ts +++ b/ui/src/pages/chat/realtime-talk-google-live.ts @@ -69,7 +69,7 @@ function googleLiveVideoMessage(frame: RealtimeTalkVideoFrame): unknown { }, }; } -const GOOGLE_LIVE_WEBSOCKET_OPEN_TIMEOUT_MS = 30_000; +const GOOGLE_LIVE_SETUP_TIMEOUT_MS = 30_000; // Browser sessions can still pin a 2.5 model, whose text and tool-response wire // contract differs from the 3.1 default carried in new session metadata. @@ -107,7 +107,7 @@ function buildGoogleLiveUrl(session: RealtimeTalkJsonPcmWebSocketSessionResult): export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { private ws: WebSocket | null = null; - private openTimeout: ReturnType | null = null; + private setupTimeout: ReturnType | null = null; private media: MediaStream | null = null; private inputContext: AudioContext | null = null; private outputContext: AudioContext | null = null; @@ -186,19 +186,18 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { const ws = new WebSocket(wsUrl); this.ws = ws; ws.binaryType = "arraybuffer"; - this.openTimeout = globalThis.setTimeout(() => { + this.setupTimeout = globalThis.setTimeout(() => { if (this.closed || this.ws !== ws) { return; } - this.openTimeout = null; + this.setupTimeout = null; this.ctx.callbacks.onStatus?.( "error", - `Realtime connection timed out after ${GOOGLE_LIVE_WEBSOCKET_OPEN_TIMEOUT_MS}ms`, + `Realtime connection timed out after ${GOOGLE_LIVE_SETUP_TIMEOUT_MS}ms`, ); this.stop(); - }, GOOGLE_LIVE_WEBSOCKET_OPEN_TIMEOUT_MS); + }, GOOGLE_LIVE_SETUP_TIMEOUT_MS); ws.addEventListener("open", () => { - this.clearOpenTimeout(ws); if (this.closed || this.ws !== ws) { return; } @@ -206,7 +205,7 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { this.startMicrophonePump(); }); ws.addEventListener("message", (event) => { - void this.handleMessage(event.data); + void this.handleMessage(ws, event.data); }); ws.addEventListener("close", () => { this.failSocket(ws, "Realtime connection closed"); @@ -232,7 +231,7 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { this.mediaSetupController?.abort(); this.mediaSetupController = null; this.setupComplete = false; - this.clearOpenTimeout(); + this.clearSetupTimeout(); for (const controller of this.consultAbortControllers) { controller.abort(); } @@ -253,15 +252,15 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { this.ws = null; } - private clearOpenTimeout(ws?: WebSocket): void { + private clearSetupTimeout(ws?: WebSocket): void { // A stopped socket can emit after a replacement starts. Keep stale events // from clearing the replacement attempt's startup deadline. if (ws && this.ws !== ws) { return; } - if (this.openTimeout !== null) { - globalThis.clearTimeout(this.openTimeout); - this.openTimeout = null; + if (this.setupTimeout !== null) { + globalThis.clearTimeout(this.setupTimeout); + this.setupTimeout = null; } } @@ -271,7 +270,7 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { if (this.closed || this.ws !== ws) { return; } - this.clearOpenTimeout(ws); + this.clearSetupTimeout(ws); this.ctx.callbacks.onStatus?.("error", message); this.stop(); } @@ -304,8 +303,8 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { return false; } - private async handleMessage(data: unknown): Promise { - if (this.closed) { + private async handleMessage(ws: WebSocket, data: unknown): Promise { + if (this.closed || this.ws !== ws) { return; } let message: GoogleLiveMessage; @@ -314,11 +313,12 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { } catch { return; } - if (this.closed) { + if (this.closed || this.ws !== ws) { return; } if (message.setupComplete) { this.setupComplete = true; + this.clearSetupTimeout(ws); this.ctx.callbacks.onStatus?.("listening"); this.emitTalkEvent({ type: "session.ready" }); this.startVideoFrames(); From 2cdc14e0c06497db8425c61eb046947362993b14 Mon Sep 17 00:00:00 2001 From: Alix-007 Date: Fri, 17 Jul 2026 20:57:14 +0800 Subject: [PATCH 3/9] test(ui): cover video setup timeout cleanup --- .../realtime-talk-google-live-video.test.ts | 53 +++++++++++++++++++ .../pages/chat/realtime-talk-google-live.ts | 1 + 2 files changed, 54 insertions(+) diff --git a/ui/src/pages/chat/realtime-talk-google-live-video.test.ts b/ui/src/pages/chat/realtime-talk-google-live-video.test.ts index 563034745e7f..7a3a02f01112 100644 --- a/ui/src/pages/chat/realtime-talk-google-live-video.test.ts +++ b/ui/src/pages/chat/realtime-talk-google-live-video.test.ts @@ -370,4 +370,57 @@ describe("Google Live Video Talk", () => { transport.stop(); }); + + it("releases camera media when Google setup times out", async () => { + const audioStop = vi.fn(); + const videoStop = vi.fn(); + const audioTrack = { stop: audioStop } as unknown as MediaStreamTrack; + const videoTrack = Object.assign(new EventTarget(), { + stop: videoStop, + readyState: "live", + enabled: true, + muted: false, + }) as unknown as MediaStreamTrack; + const audio = { + getAudioTracks: () => [audioTrack], + getTracks: () => [audioTrack], + } as unknown as MediaStream; + const camera = { + getVideoTracks: () => [videoTrack], + getTracks: () => [videoTrack], + } as unknown as MediaStream; + const getUserMedia = vi.fn().mockResolvedValueOnce(audio).mockResolvedValueOnce(camera); + vi.stubGlobal("navigator", { mediaDevices: { getUserMedia } }); + const originalCreateElement = document.createElement.bind(document); + vi.spyOn(document, "createElement").mockImplementation((tagName: string) => { + const element = originalCreateElement(tagName); + if (element instanceof HTMLVideoElement) { + vi.spyOn(element, "play").mockResolvedValue(undefined); + } + return element; + }); + const onStatus = vi.fn(); + const onVideoStream = vi.fn(); + const transport = createTransport({ onStatus, onVideoStream }); + + await transport.start(); + await transport.setVideoEnabled(true); + const ws = FakeGoogleLiveWebSocket.instance; + if (!ws) { + throw new Error("missing Google Live WebSocket"); + } + ws.emitOpen(); + await vi.advanceTimersByTimeAsync(30_000); + + expect(onStatus).toHaveBeenCalledExactlyOnceWith( + "error", + "Realtime connection timed out after 30000ms", + ); + expect(audioStop).toHaveBeenCalledOnce(); + expect(videoStop).toHaveBeenCalledOnce(); + expect(getUserMedia).toHaveBeenNthCalledWith(2, { video: true }); + expect(onVideoStream).toHaveBeenNthCalledWith(1, camera); + expect(onVideoStream).toHaveBeenLastCalledWith(null); + expect(ws.readyState).toBe(3); + }); }); diff --git a/ui/src/pages/chat/realtime-talk-google-live.ts b/ui/src/pages/chat/realtime-talk-google-live.ts index 428c4f3d10d5..06401349c462 100644 --- a/ui/src/pages/chat/realtime-talk-google-live.ts +++ b/ui/src/pages/chat/realtime-talk-google-live.ts @@ -69,6 +69,7 @@ function googleLiveVideoMessage(frame: RealtimeTalkVideoFrame): unknown { }, }; } + const GOOGLE_LIVE_SETUP_TIMEOUT_MS = 30_000; // Browser sessions can still pin a 2.5 model, whose text and tool-response wire From c5372430b3e2812c46849a3cc5f309e92eeadcee Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 17:20:17 +0800 Subject: [PATCH 4/9] fix(talk): gate Google Live activation on setup --- .../realtime-talk-google-live-lifecycle.ts | 145 +++++++++++++++ .../pages/chat/realtime-talk-google-live.ts | 170 ++++++++++++------ 2 files changed, 257 insertions(+), 58 deletions(-) create mode 100644 ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts diff --git a/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts b/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts new file mode 100644 index 000000000000..1d12604af401 --- /dev/null +++ b/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts @@ -0,0 +1,145 @@ +import type { + RealtimeTalkJsonPcmWebSocketSessionResult, + RealtimeTalkTransportStartResult, +} from "./realtime-talk-shared.ts"; + +const GOOGLE_LIVE_WEBSOCKET_HOST = "generativelanguage.googleapis.com"; +const GOOGLE_LIVE_WEBSOCKET_PATH = + /^\/ws\/google\.ai\.generativelanguage\.v[0-9a-z]+\.GenerativeService\.BidiGenerateContent(?:Constrained)?$/; +export const GOOGLE_LIVE_SETUP_TIMEOUT_MS = 30_000; + +export function buildGoogleLiveUrl(session: RealtimeTalkJsonPcmWebSocketSessionResult): string { + let url: URL; + try { + url = new URL(session.websocketUrl); + } catch { + throw new Error("Invalid Google Live WebSocket URL"); + } + if (url.protocol !== "wss:") { + throw new Error("Google Live WebSocket URL must use wss://"); + } + if (url.hostname.toLowerCase() !== GOOGLE_LIVE_WEBSOCKET_HOST) { + throw new Error("Untrusted Google Live WebSocket host"); + } + if (url.username || url.password) { + throw new Error("Google Live WebSocket URL must not include credentials"); + } + if (!GOOGLE_LIVE_WEBSOCKET_PATH.test(url.pathname)) { + throw new Error("Untrusted Google Live WebSocket path"); + } + url.search = ""; + url.searchParams.set("access_token", session.clientSecret); + return url.toString(); +} + +export type GoogleLiveConnectionState = + | "idle" + | "connecting" + | "ready" + | "active" + | "cancelled" + | "failed"; + +type StartupWaiter = { + resolve: (result: RealtimeTalkTransportStartResult) => void; + reject: (error: Error) => void; +}; + +export class GoogleLiveConnectionLifecycle { + private state: GoogleLiveConnectionState = "idle"; + private socket: WebSocket | null = null; + private waiter: StartupWaiter | null = null; + private error: Error | null = null; + + get currentState(): GoogleLiveConnectionState { + return this.state; + } + + get isActive(): boolean { + return this.state === "active"; + } + + get setupComplete(): boolean { + return this.state === "ready" || this.state === "active"; + } + + begin(socket: WebSocket): Promise { + this.state = "connecting"; + this.socket = socket; + this.error = null; + return new Promise((resolve, reject) => { + this.waiter = { resolve, reject }; + }); + } + + markReady(socket: WebSocket): boolean { + if (this.state !== "connecting" || this.socket !== socket) { + return false; + } + this.state = "ready"; + this.takeWaiter()?.resolve("ready"); + return true; + } + + finishStart(result: RealtimeTalkTransportStartResult): RealtimeTalkTransportStartResult { + if (this.error) { + throw this.error; + } + return this.state === "cancelled" ? "cancelled" : result; + } + + activate(): boolean { + if (this.state === "active") { + return false; + } + if (this.error) { + throw this.error; + } + if (this.state === "cancelled" || this.state === "idle") { + return false; + } + if (this.state !== "ready") { + throw new Error("Google Live transport activated before setup completed"); + } + this.state = "active"; + return true; + } + + failStartup(socket: WebSocket, error: Error): boolean { + if (this.socket !== socket || (this.state !== "connecting" && this.state !== "ready")) { + return false; + } + this.state = "failed"; + this.error = error; + this.takeWaiter()?.reject(error); + return true; + } + + cancel(): void { + if (this.state === "cancelled" || this.state === "failed") { + return; + } + this.state = "cancelled"; + this.takeWaiter()?.resolve("cancelled"); + } + + private takeWaiter(): StartupWaiter | null { + const waiter = this.waiter; + this.waiter = null; + return waiter; + } +} + +export function runRealtimeTalkCleanup(steps: Array<() => void>): void { + let firstError: unknown; + for (const step of steps) { + try { + step(); + } catch (error) { + firstError ??= error; + } + } + if (firstError) { + throw firstError; + } +} diff --git a/ui/src/pages/chat/realtime-talk-google-live.ts b/ui/src/pages/chat/realtime-talk-google-live.ts index d7fabdb2eb93..38160a287524 100644 --- a/ui/src/pages/chat/realtime-talk-google-live.ts +++ b/ui/src/pages/chat/realtime-talk-google-live.ts @@ -9,6 +9,12 @@ import { RealtimeTalkPcmOutputQueue, } from "./realtime-talk-audio.ts"; import { RealtimeTalkCameraController } from "./realtime-talk-camera-controller.ts"; +import { + buildGoogleLiveUrl, + GoogleLiveConnectionLifecycle, + GOOGLE_LIVE_SETUP_TIMEOUT_MS, + runRealtimeTalkCleanup, +} from "./realtime-talk-google-live-lifecycle.ts"; import { openRealtimeTalkCamera, openRealtimeTalkInput } from "./realtime-talk-input.ts"; import type { RealtimeTalkJsonPcmWebSocketSessionResult } from "./realtime-talk-shared.ts"; import { @@ -57,9 +63,6 @@ type PendingFunctionCall = { args: unknown; }; -const GOOGLE_LIVE_WEBSOCKET_HOST = "generativelanguage.googleapis.com"; -const GOOGLE_LIVE_WEBSOCKET_PATH = - /^\/ws\/google\.ai\.generativelanguage\.v[0-9a-z]+\.GenerativeService\.BidiGenerateContent(?:Constrained)?$/; const GOOGLE_LIVE_VIDEO_FRAME_INTERVAL_MS = 1_000; const GOOGLE_LIVE_VIDEO_MESSAGE_MAX_BYTES = 512 * 1024; @@ -71,8 +74,6 @@ function googleLiveVideoMessage(frame: RealtimeTalkVideoFrame): unknown { }; } -const GOOGLE_LIVE_SETUP_TIMEOUT_MS = 30_000; - // Browser sessions can still pin a 2.5 model, whose text and tool-response wire // contract differs from the 3.1 default carried in new session metadata. function isGemini31LiveModel(model: string | undefined): boolean { @@ -83,30 +84,6 @@ function isGemini31LiveModel(model: string | undefined): boolean { return modelId.startsWith("gemini-3.1-") && modelId.includes("-live"); } -function buildGoogleLiveUrl(session: RealtimeTalkJsonPcmWebSocketSessionResult): string { - let url: URL; - try { - url = new URL(session.websocketUrl); - } catch { - throw new Error("Invalid Google Live WebSocket URL"); - } - if (url.protocol !== "wss:") { - throw new Error("Google Live WebSocket URL must use wss://"); - } - if (url.hostname.toLowerCase() !== GOOGLE_LIVE_WEBSOCKET_HOST) { - throw new Error("Untrusted Google Live WebSocket host"); - } - if (url.username || url.password) { - throw new Error("Google Live WebSocket URL must not include credentials"); - } - if (!GOOGLE_LIVE_WEBSOCKET_PATH.test(url.pathname)) { - throw new Error("Untrusted Google Live WebSocket path"); - } - url.search = ""; - url.searchParams.set("access_token", session.clientSecret); - return url.toString(); -} - export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { private ws: WebSocket | null = null; private setupTimeout: ReturnType | null = null; @@ -118,7 +95,8 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { private closed = false; private mediaSetupController: AbortController | null = null; private readonly camera: RealtimeTalkCameraController; - private setupComplete = false; + private readonly lifecycle = new GoogleLiveConnectionLifecycle(); + private cameraPublished = false; private videoFramesActive = false; private hasSentVideoFrame = false; private videoFrameTimer: ReturnType | null = null; @@ -137,9 +115,20 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { getDeviceId: () => this.ctx.videoDeviceId, setDeviceId: (deviceId) => (this.ctx.videoDeviceId = deviceId), isClosed: () => this.closed, - onStream: (stream) => this.ctx.callbacks.onVideoStream?.(stream), + onStream: (stream) => { + if (stream) { + if (!this.lifecycle.isActive) { + return; + } + this.cameraPublished = true; + this.ctx.callbacks.onVideoStream?.(stream); + } else if (this.cameraPublished) { + this.cameraPublished = false; + this.ctx.callbacks.onVideoStream?.(null); + } + }, onAcquired: () => { - if (this.setupComplete) { + if (this.lifecycle.isActive) { this.startVideoFrames(); } }, @@ -156,6 +145,7 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { } const wsUrl = buildGoogleLiveUrl(this.session); this.closed = false; + this.cameraPublished = false; this.mediaSetupController?.abort(); const mediaSetupController = new AbortController(); this.mediaSetupController = mediaSetupController; @@ -181,13 +171,10 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { this.media = media; this.inputContext = new AudioContext({ sampleRate: this.session.audio.inputSampleRateHz }); this.outputContext = new AudioContext({ sampleRate: this.session.audio.outputSampleRateHz }); - if (this.ctx.callbacks.onInputLevel) { - this.inputMeter = new RealtimeTalkMediaStreamMeter(this.ctx.callbacks.onInputLevel); - this.inputMeter.start(this.media, this.inputContext); - } const ws = new WebSocket(wsUrl); this.ws = ws; ws.binaryType = "arraybuffer"; + const startup = this.lifecycle.begin(ws); this.setupTimeout = globalThis.setTimeout(() => { if (this.closed || this.ws !== ws) { return; @@ -203,7 +190,6 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { return; } this.send(this.session.initialMessage ?? { setup: {} }); - this.startMicrophonePump(); }); ws.addEventListener("message", (event) => { void this.handleMessage(ws, event.data); @@ -214,7 +200,46 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { ws.addEventListener("error", () => { this.failConnection(ws, "Realtime connection failed"); }); - return "ready"; + return this.lifecycle.finishStart(await startup); + } + + activate(): void { + if (this.closed || !this.lifecycle.activate()) { + return; + } + try { + this.ctx.callbacks.onStatus?.("listening"); + if (this.closed) { + return; + } + this.emitTalkEvent({ type: "session.ready" }); + if (this.closed) { + return; + } + if (this.ctx.callbacks.onInputLevel && this.media && this.inputContext) { + this.inputMeter = new RealtimeTalkMediaStreamMeter(this.ctx.callbacks.onInputLevel); + this.inputMeter.start(this.media, this.inputContext); + } + if (this.closed) { + return; + } + this.startMicrophonePump(); + if (this.camera.stream && !this.cameraPublished) { + this.cameraPublished = true; + this.ctx.callbacks.onVideoStream?.(this.camera.stream); + } + if (this.closed) { + return; + } + this.startVideoFrames(); + } catch (error) { + try { + this.stop({ emitClosed: false }); + } catch { + // Preserve the activation callback as the terminal cause after cleanup. + } + throw error; + } } async setVideoEnabled(enabled: boolean): Promise { @@ -226,40 +251,60 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { } stop(options?: { emitClosed?: boolean }): void { - const emitClosed = !this.closed && options?.emitClosed !== false; + const emitClosed = !this.closed && this.lifecycle.isActive && options?.emitClosed !== false; this.closed = true; + this.lifecycle.cancel(); + let firstError: unknown; try { if (emitClosed) { this.emitTalkEvent({ type: "session.closed", final: true }); } - } finally { + } catch (error) { + firstError = error; + } + try { this.releaseResources(); + } catch (error) { + firstError ??= error; + } + if (firstError) { + throw firstError; } } private releaseResources(): void { - this.mediaSetupController?.abort(); + const mediaSetupController = this.mediaSetupController; this.mediaSetupController = null; - this.setupComplete = false; this.clearSetupTimeout(); - for (const controller of this.consultAbortControllers) { - controller.abort(); - } + const consultAbortControllers = [...this.consultAbortControllers]; this.consultAbortControllers.clear(); this.pendingCalls.clear(); - this.inputPump.stop(); - this.inputMeter?.stop(); + const inputMeter = this.inputMeter; this.inputMeter = null; - this.media?.getTracks().forEach((track) => track.stop()); + const media = this.media; this.media = null; - this.camera.release(); - this.stopOutput(); - void this.inputContext?.close(); + const inputContext = this.inputContext; this.inputContext = null; - void this.outputContext?.close(); + const outputContext = this.outputContext; this.outputContext = null; - this.ws?.close(); + const ws = this.ws; this.ws = null; + runRealtimeTalkCleanup([ + () => mediaSetupController?.abort(), + ...consultAbortControllers.map((controller) => () => controller.abort()), + () => this.inputPump.stop(), + () => inputMeter?.stop(), + ...(media?.getTracks() ?? []).map((track) => () => track.stop()), + () => this.camera.release(), + () => this.stopOutput(), + () => { + void inputContext?.close(); + }, + () => { + void outputContext?.close(); + }, + () => ws?.close(), + ]); } private clearSetupTimeout(): void { @@ -275,6 +320,14 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { if (this.closed || this.ws !== ws) { return; } + if (this.lifecycle.failStartup(ws, new Error(detail))) { + try { + this.stop({ emitClosed: false }); + } catch { + // Startup rejection owns terminal precedence; cleanup still ran to completion. + } + return; + } try { this.ctx.callbacks.onStatus?.("error", detail); } finally { @@ -324,12 +377,13 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { if (this.closed || this.ws !== ws) { return; } - if (message.setupComplete) { - this.setupComplete = true; + if (message.setupComplete && this.lifecycle.markReady(ws)) { this.clearSetupTimeout(); - this.ctx.callbacks.onStatus?.("listening"); - this.emitTalkEvent({ type: "session.ready" }); - this.startVideoFrames(); + } + // The parent session adopts the candidate after start() resolves. Provider + // events remain provisional until activate() publishes that ownership. + if (!this.lifecycle.isActive) { + return; } const content = message.serverContent; if (content?.interrupted) { From d4a477908cd9541fcdf47d33a3c86d966e6fec54 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 17:20:20 +0800 Subject: [PATCH 5/9] test(talk): cover Google Live startup transitions --- .../realtime-talk-google-live-timeout.test.ts | 149 ++++++++++++----- .../chat/realtime-talk-google-live.test.ts | 154 +++++++++++------- 2 files changed, 208 insertions(+), 95 deletions(-) diff --git a/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts b/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts index 5eba5aabcda3..e3d1d62d4b5a 100644 --- a/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts +++ b/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts @@ -35,6 +35,14 @@ class FakeGoogleLiveWebSocket extends EventTarget { emitMessage(message: unknown): void { this.dispatchEvent(new MessageEvent("message", { data: JSON.stringify(message) })); } + + emitClose(): void { + this.dispatchEvent(new Event("close")); + } + + emitError(): void { + this.dispatchEvent(new Event("error")); + } } class FakeAudioContext { @@ -94,6 +102,15 @@ function latestSocket(): FakeGoogleLiveWebSocket { return socket; } +async function beginTransport(transport: GoogleLiveRealtimeTalkTransport): Promise<{ + start: Promise<"ready" | "cancelled">; + socket: FakeGoogleLiveWebSocket; +}> { + const start = transport.start(); + await vi.advanceTimersByTimeAsync(0); + return { start, socket: latestSocket() }; +} + describe("Google Live setup timeout", () => { beforeEach(() => { vi.useFakeTimers(); @@ -121,22 +138,18 @@ describe("Google Live setup timeout", () => { const onTalkEvent = vi.fn(); const transport = createTransport({ onStatus, onTalkEvent }); - await expect(transport.start()).resolves.toBe("ready"); - const socket = latestSocket(); + const { start, socket } = await beginTransport(transport); socket.readyState = 0; + const rejected = expect(start).rejects.toThrow("Realtime connection timed out after 30000ms"); await vi.advanceTimersByTimeAsync(SETUP_TIMEOUT_MS); - expect(onStatus).toHaveBeenCalledExactlyOnceWith( - "error", - "Realtime connection timed out after 30000ms", - ); + await rejected; + expect(onStatus).not.toHaveBeenCalled(); expect(stopInputTrack).toHaveBeenCalledOnce(); expect(audioContexts).toHaveLength(2); expect(audioContexts.every((context) => context.close.mock.calls.length === 1)).toBe(true); expect(socket.readyState).toBe(3); - expect(onTalkEvent).toHaveBeenCalledExactlyOnceWith( - expect.objectContaining({ type: "session.closed", final: true }), - ); + expect(onTalkEvent).not.toHaveBeenCalled(); expect(vi.getTimerCount()).toBe(0); }); @@ -144,51 +157,106 @@ describe("Google Live setup timeout", () => { const onStatus = vi.fn(); const transport = createTransport({ onStatus }); - await transport.start(); - latestSocket().emitOpen(); + const { start, socket } = await beginTransport(transport); + socket.emitOpen(); + const rejected = expect(start).rejects.toThrow("Realtime connection timed out after 30000ms"); await vi.advanceTimersByTimeAsync(SETUP_TIMEOUT_MS); - expect(onStatus).toHaveBeenCalledExactlyOnceWith( - "error", - "Realtime connection timed out after 30000ms", - ); + await rejected; + expect(onStatus).not.toHaveBeenCalled(); expect(stopInputTrack).toHaveBeenCalledOnce(); expect(audioContexts.every((context) => context.close.mock.calls.length === 1)).toBe(true); }); - it.each(["status", "talk event"] as const)( - "releases timed-out resources when the terminal %s callback throws", - async (callbackKind) => { - const throwingCallback = vi.fn(() => { - throw new Error("consumer failed"); - }); - const transport = createTransport( - callbackKind === "status" - ? { onStatus: throwingCallback } - : { onTalkEvent: throwingCallback }, - ); + it("does not publish provisional terminal callbacks when setup times out", async () => { + const onStatus = vi.fn(() => { + throw new Error("status callback must remain provisional"); + }); + const onTalkEvent = vi.fn(() => { + throw new Error("talk callback must remain provisional"); + }); + const transport = createTransport({ onStatus, onTalkEvent }); - await transport.start(); - const socket = latestSocket(); - socket.readyState = 0; - expect(() => vi.advanceTimersByTime(SETUP_TIMEOUT_MS)).toThrow("consumer failed"); + const { start, socket } = await beginTransport(transport); + socket.readyState = 0; + const rejected = expect(start).rejects.toThrow("Realtime connection timed out after 30000ms"); + await vi.advanceTimersByTimeAsync(SETUP_TIMEOUT_MS); - expect(stopInputTrack).toHaveBeenCalledOnce(); - expect(audioContexts.every((context) => context.close.mock.calls.length === 1)).toBe(true); - expect(socket.readyState).toBe(3); - expect(vi.getTimerCount()).toBe(0); - }, - ); + await rejected; + expect(onStatus).not.toHaveBeenCalled(); + expect(onTalkEvent).not.toHaveBeenCalled(); + expect(stopInputTrack).toHaveBeenCalledOnce(); + expect(audioContexts.every((context) => context.close.mock.calls.length === 1)).toBe(true); + expect(socket.readyState).toBe(3); + expect(vi.getTimerCount()).toBe(0); + }); + + it.each([ + ["close", "Realtime connection closed"], + ["error", "Realtime connection failed"], + ] as const)("rejects startup when the WebSocket emits %s", async (event, detail) => { + const onStatus = vi.fn(); + const transport = createTransport({ onStatus }); + const { start, socket } = await beginTransport(transport); + const rejected = expect(start).rejects.toThrow(detail); + + if (event === "close") { + socket.emitClose(); + } else { + socket.emitError(); + } + + await rejected; + expect(onStatus).not.toHaveBeenCalled(); + expect(stopInputTrack).toHaveBeenCalledOnce(); + expect(socket.readyState).toBe(3); + expect(vi.getTimerCount()).toBe(0); + }); + + it("rejects when the socket closes after setup but before activation", async () => { + const onStatus = vi.fn(); + const transport = createTransport({ onStatus }); + const { start, socket } = await beginTransport(transport); + socket.emitOpen(); + socket.emitMessage({ setupComplete: {} }); + await Promise.resolve(); + const rejected = expect(start).rejects.toThrow("Realtime connection closed"); + + socket.emitClose(); + + await rejected; + expect(onStatus).not.toHaveBeenCalled(); + expect(stopInputTrack).toHaveBeenCalledOnce(); + expect(socket.readyState).toBe(3); + expect(vi.getTimerCount()).toBe(0); + }); + + it("releases resources when a readiness callback throws during activation", async () => { + const onStatus = vi.fn(() => { + throw new Error("consumer failed"); + }); + const transport = createTransport({ onStatus }); + const { start, socket } = await beginTransport(transport); + socket.emitOpen(); + socket.emitMessage({ setupComplete: {} }); + await expect(start).resolves.toBe("ready"); + + expect(() => transport.activate()).toThrow("consumer failed"); + expect(stopInputTrack).toHaveBeenCalledOnce(); + expect(socket.readyState).toBe(3); + expect(audioContexts.every((context) => context.close.mock.calls.length === 1)).toBe(true); + }); it("clears the deadline after Google setup completes", async () => { const onStatus = vi.fn(); const transport = createTransport({ onStatus }); - await transport.start(); - const socket = latestSocket(); + const { start, socket } = await beginTransport(transport); socket.emitOpen(); socket.emitMessage({ setupComplete: {} }); - await Promise.resolve(); + await expect(start).resolves.toBe("ready"); + expect(onStatus).not.toHaveBeenCalled(); + transport.activate(); await vi.advanceTimersByTimeAsync(SETUP_TIMEOUT_MS); expect(onStatus).toHaveBeenCalledExactlyOnceWith("listening"); @@ -199,8 +267,9 @@ describe("Google Live setup timeout", () => { const onStatus = vi.fn(); const transport = createTransport({ onStatus }); - await transport.start(); + const { start } = await beginTransport(transport); transport.stop(); + await expect(start).resolves.toBe("cancelled"); await vi.advanceTimersByTimeAsync(SETUP_TIMEOUT_MS); expect(onStatus).not.toHaveBeenCalled(); diff --git a/ui/src/pages/chat/realtime-talk-google-live.test.ts b/ui/src/pages/chat/realtime-talk-google-live.test.ts index e59a073caf2e..f69e9c2cb7f5 100644 --- a/ui/src/pages/chat/realtime-talk-google-live.test.ts +++ b/ui/src/pages/chat/realtime-talk-google-live.test.ts @@ -223,6 +223,26 @@ function latestWebSocket(): MockGoogleLiveWebSocket { return ws; } +async function beginTransport(transport: GoogleLiveRealtimeTalkTransport): Promise<{ + start: Promise<"ready" | "cancelled">; + ws: MockGoogleLiveWebSocket; +}> { + const start = transport.start(); + await waitForFast(() => expect(wsInstances).toHaveLength(1)); + return { start, ws: latestWebSocket() }; +} + +async function startTransport( + transport: GoogleLiveRealtimeTalkTransport, +): Promise { + const { start, ws } = await beginTransport(transport); + ws.emitOpen(); + ws.emitMessage(encodeJsonFrame({ setupComplete: {} })); + await expect(start).resolves.toBe("ready"); + transport.activate(); + return ws; +} + function pumpMicrophone(samples: Float32Array): void { const processor = inputProcessors.at(-1); if (!processor) { @@ -275,12 +295,13 @@ describe("GoogleLiveRealtimeTalkTransport", () => { { callbacks: {}, client: createClient(), sessionKey: "main" }, ); - await transport.start(); + const { start } = await beginTransport(transport); expect(latestWebSocket().url).toBe( "wss://generativelanguage.googleapis.com/ws/google.ai.generativelanguage.v1alpha.GenerativeService.BidiGenerateContentConstrained?access_token=auth_tokens%2Fbrowser-session", ); transport.stop(); + await expect(start).resolves.toBe("cancelled"); }); it.each([ @@ -301,7 +322,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { it("captures from the selected microphone with an exact constraint", async () => { const transport = createTransport({}, createClient(), "usb-mic"); - await transport.start(); + const { start } = await beginTransport(transport); expect(getUserMedia).toHaveBeenCalledWith({ audio: { @@ -312,13 +333,13 @@ describe("GoogleLiveRealtimeTalkTransport", () => { }, }); transport.stop(); + await expect(start).resolves.toBe("cancelled"); }); it("keeps the microphone processor inaudible locally", async () => { const transport = createTransport(); - await transport.start(); - latestWebSocket().emitOpen(); + await startTransport(transport); const processor = inputProcessors.at(-1); const sink = inputSinks.at(-1); @@ -359,12 +380,18 @@ describe("GoogleLiveRealtimeTalkTransport", () => { const onTalkEvent = vi.fn(); const transport = createTransport({ onStatus, onTalkEvent }); - await transport.start(); - const ws = latestWebSocket(); + const { start, ws } = await beginTransport(transport); + ws.emitOpen(); ws.emitMessage(encodeJsonFrame({ setupComplete: {} })); + await expect(start).resolves.toBe("ready"); expect(ws.binaryType).toBe("arraybuffer"); - await waitForFast(() => expect(onStatus).toHaveBeenCalledWith("listening")); + expect(onStatus).not.toHaveBeenCalled(); + expect(onTalkEvent).not.toHaveBeenCalled(); + transport.activate(); + transport.activate(); + expect(onStatus).toHaveBeenCalledWith("listening"); + expect(onStatus).toHaveBeenCalledOnce(); const readyEvent = requireFirstTalkEvent(onTalkEvent); expect(readyEvent.type).toBe("session.ready"); expect(readyEvent.sessionId).toBe("main:google:provider-websocket"); @@ -376,9 +403,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { const onTranscript = vi.fn(); const transport = createTransport({ onStatus, onTranscript }); - await expect(transport.start()).resolves.toBe("ready"); - const ws = latestWebSocket(); - ws.emitOpen(); + const ws = await startTransport(transport); pumpMicrophone(new Float32Array(4096)); ws.emitClose(); @@ -406,9 +431,8 @@ describe("GoogleLiveRealtimeTalkTransport", () => { const onTranscript = vi.fn(); const transport = createTransport({ onStatus, onTranscript }); - await expect(transport.start()).resolves.toBe("ready"); - const ws = latestWebSocket(); - ws.emitOpen(); + const ws = await startTransport(transport); + onStatus.mockClear(); ws.emitError(); ws.emitClose(); @@ -429,18 +453,25 @@ describe("GoogleLiveRealtimeTalkTransport", () => { it.each(["status", "talk event"] as const)( "releases socket resources when the terminal %s callback throws", async (callbackKind) => { - const throwingCallback = vi.fn(() => { - throw new Error("consumer failed"); - }); const transport = createTransport( callbackKind === "status" - ? { onStatus: throwingCallback } - : { onTalkEvent: throwingCallback }, + ? { + onStatus: vi.fn((status) => { + if (status === "error") { + throw new Error("consumer failed"); + } + }), + } + : { + onTalkEvent: vi.fn((event) => { + if (event.type === "session.closed") { + throw new Error("consumer failed"); + } + }), + }, ); - await expect(transport.start()).resolves.toBe("ready"); - const ws = latestWebSocket(); - ws.emitOpen(); + const ws = await startTransport(transport); expect(() => ws.emitError()).toThrow("consumer failed"); expect(stopInputTrack).toHaveBeenCalledOnce(); @@ -452,12 +483,30 @@ describe("GoogleLiveRealtimeTalkTransport", () => { }, ); + it("finishes cleanup when the input-level callback throws during activation", async () => { + const onInputLevel = vi.fn(() => { + throw new Error("meter callback failed"); + }); + const transport = createTransport({ onInputLevel }); + const { start, ws } = await beginTransport(transport); + ws.emitOpen(); + ws.emitMessage(encodeJsonFrame({ setupComplete: {} })); + await expect(start).resolves.toBe("ready"); + + expect(() => transport.activate()).toThrow("meter callback failed"); + expect(stopInputTrack).toHaveBeenCalledOnce(); + expect(inputProcessors).toHaveLength(0); + expect(ws.readyState).toBe(3); + for (const context of audioContexts) { + expect(context.close).toHaveBeenCalledOnce(); + } + }); + it("reports microphone activity and resets it when stopped", async () => { const onInputLevel = vi.fn(); const transport = createTransport({ onInputLevel }); - await transport.start(); - latestWebSocket().emitOpen(); + await startTransport(transport); pumpMicrophone(new Float32Array(4096)); pumpMicrophone(new Float32Array(4096).fill(0.25)); transport.stop(); @@ -470,8 +519,11 @@ describe("GoogleLiveRealtimeTalkTransport", () => { const onStatus = vi.fn(); const transport = createTransport({ onStatus }); - await transport.start(); - latestWebSocket().emitMessage(new Blob([JSON.stringify({ setupComplete: {} })])); + const { start, ws } = await beginTransport(transport); + ws.emitOpen(); + ws.emitMessage(new Blob([JSON.stringify({ setupComplete: {} })])); + await expect(start).resolves.toBe("ready"); + transport.activate(); await waitForFast(() => expect(onStatus).toHaveBeenCalledWith("listening")); }); @@ -479,8 +531,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { it("stops queued output when Google Live sends interruption", async () => { const onTalkEvent = vi.fn(); const transport = createTransport({ onTalkEvent }); - await transport.start(); - const ws = latestWebSocket(); + const ws = await startTransport(transport); ws.emitMessage( encodeJsonFrame({ @@ -508,8 +559,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { const onStatus = vi.fn(); const onTalkEvent = vi.fn(); const transport = createTransport({ onStatus, onTalkEvent }); - await transport.start(); - const ws = latestWebSocket(); + const ws = await startTransport(transport); ws.emitMessage( encodeJsonFrame({ @@ -557,8 +607,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { it("rejects an oversized first frame before decoding provider audio", async () => { const onStatus = vi.fn(); const transport = createTransport({ onStatus }); - await transport.start(); - const ws = latestWebSocket(); + const ws = await startTransport(transport); ws.emitMessage( encodeJsonFrame({ @@ -592,8 +641,9 @@ describe("GoogleLiveRealtimeTalkTransport", () => { const onTalkEvent = vi.fn(); const transport = createTransport({ onTalkEvent, onTranscript }); - await transport.start(); - latestWebSocket().emitMessage( + const ws = await startTransport(transport); + onTalkEvent.mockClear(); + ws.emitMessage( encodeJsonFrame({ serverContent: { inputTranscription: { text: "hello", finished: true }, @@ -638,8 +688,9 @@ describe("GoogleLiveRealtimeTalkTransport", () => { const onTranscript = vi.fn(() => transport.stop()); const transport = createTransport({ onTalkEvent, onTranscript }); - await transport.start(); - latestWebSocket().emitMessage( + const ws = await startTransport(transport); + onTalkEvent.mockClear(); + ws.emitMessage( encodeJsonFrame({ serverContent: { inputTranscription: { text: "overflow", finished: true }, @@ -657,22 +708,22 @@ describe("GoogleLiveRealtimeTalkTransport", () => { it("silently disposes a provisional Google Live transport", async () => { const onTalkEvent = vi.fn(); const transport = createTransport({ onTalkEvent }); - await transport.start(); - onTalkEvent.mockClear(); + const { start, ws } = await beginTransport(transport); transport.stop({ emitClosed: false }); + await expect(start).resolves.toBe("cancelled"); expect(onTalkEvent).not.toHaveBeenCalled(); - expect(latestWebSocket().readyState).toBe(3); + expect(ws.readyState).toBe(3); }); it("ignores late WebSocket events after stop", async () => { const onStatus = vi.fn(); const transport = createTransport({ onStatus }); - await transport.start(); - const ws = latestWebSocket(); + const { start, ws } = await beginTransport(transport); transport.stop(); + await expect(start).resolves.toBe("cancelled"); ws.emitOpen(); ws.emitMessage(new Blob([JSON.stringify({ setupComplete: {} })])); @@ -702,9 +753,10 @@ describe("GoogleLiveRealtimeTalkTransport", () => { }), } as unknown as RealtimeTalkTransportContext["client"]; const transport = createTransport({ onStatus }, client); - await transport.start(); + const ws = await startTransport(transport); + onStatus.mockClear(); - latestWebSocket().emitMessage( + ws.emitMessage( encodeJsonFrame({ toolCall: { functionCalls: [ @@ -744,8 +796,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { }), } as unknown as RealtimeTalkTransportContext["client"]; const transport = createTransport({}, client); - await transport.start(); - const ws = latestWebSocket(); + const ws = await startTransport(transport); ws.emitMessage( encodeJsonFrame({ @@ -802,8 +853,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { }); const transport = createTransport({ onStatus, onTalkEvent }, client); - await transport.start(); - const ws = latestWebSocket(); + const ws = await startTransport(transport); vi.spyOn(ws, "send").mockImplementation(() => { throw new Error("Google Live socket rejected the tool result"); }); @@ -868,9 +918,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { throw new Error(`unexpected request: ${method}`); }); const transport = createTransport({}, client); - await transport.start(); - const ws = latestWebSocket(); - ws.emitOpen(); + const ws = await startTransport(transport); ws.emitMessage( encodeJsonFrame({ serverContent: { @@ -941,9 +989,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { throw new Error(`unexpected request: ${method}`); }); const transport = createTransport({}, client); - await transport.start(); - const ws = latestWebSocket(); - ws.emitOpen(); + const ws = await startTransport(transport); ws.emitMessage( encodeJsonFrame({ serverContent: { @@ -1014,9 +1060,7 @@ describe("GoogleLiveRealtimeTalkTransport", () => { throw new Error(`unexpected request: ${method}`); }); const transport = createTransport({}, client); - await transport.start(); - const ws = latestWebSocket(); - ws.emitOpen(); + const ws = await startTransport(transport); ws.emitMessage( encodeJsonFrame({ serverContent: { From 476991a5239a25c9acd25a1b0df85b4fb206cae5 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 17:20:20 +0800 Subject: [PATCH 6/9] test(talk): cover Google Live video teardown --- .../realtime-talk-google-live-video.test.ts | 94 +++++++++++++++---- 1 file changed, 75 insertions(+), 19 deletions(-) diff --git a/ui/src/pages/chat/realtime-talk-google-live-video.test.ts b/ui/src/pages/chat/realtime-talk-google-live-video.test.ts index 7a3a02f01112..ab672650e708 100644 --- a/ui/src/pages/chat/realtime-talk-google-live-video.test.ts +++ b/ui/src/pages/chat/realtime-talk-google-live-video.test.ts @@ -84,6 +84,30 @@ function createTransport(callbacks: RealtimeTalkCallbacks, videoDeviceId?: strin ); } +async function beginTransport(transport: GoogleLiveRealtimeTalkTransport): Promise<{ + start: Promise<"ready" | "cancelled">; + ws: FakeGoogleLiveWebSocket; +}> { + const start = transport.start(); + await vi.advanceTimersByTimeAsync(0); + const ws = FakeGoogleLiveWebSocket.instance; + if (!ws) { + throw new Error("missing Google Live WebSocket"); + } + return { start, ws }; +} + +async function startTransport( + transport: GoogleLiveRealtimeTalkTransport, +): Promise { + const { start, ws } = await beginTransport(transport); + ws.emitOpen(); + ws.emitMessage({ setupComplete: {} }); + await expect(start).resolves.toBe("ready"); + transport.activate(); + return ws; +} + describe("Google Live Video Talk", () => { beforeEach(() => { vi.useFakeTimers(); @@ -142,16 +166,15 @@ describe("Google Live Video Talk", () => { const onVideoStream = vi.fn(); const transport = createTransport({ onStatus, onVideoStream }); - await transport.start(); + const { start, ws } = await beginTransport(transport); expect(getUserMedia).toHaveBeenCalledOnce(); expect(onVideoStream).not.toHaveBeenCalled(); await transport.setVideoEnabled(true); - const ws = FakeGoogleLiveWebSocket.instance; - if (!ws) { - throw new Error("missing Google Live WebSocket"); - } ws.emitOpen(); ws.emitMessage({ setupComplete: {} }); + await expect(start).resolves.toBe("ready"); + expect(ws.sent.some((message) => JSON.stringify(message).includes('"video"'))).toBe(false); + transport.activate(); await vi.advanceTimersByTimeAsync(0); expect(ws.sent).toContainEqual({ @@ -273,7 +296,7 @@ describe("Google Live Video Talk", () => { const onVideoStream = vi.fn(); const transport = createTransport({ onVideoStream }); - await transport.start(); + await startTransport(transport); await transport.setVideoEnabled(true); firstVideoTrack.dispatchEvent(new Event("ended")); @@ -304,7 +327,7 @@ describe("Google Live Video Talk", () => { vi.stubGlobal("navigator", { mediaDevices: { getUserMedia } }); const transport = createTransport({}); - await transport.start(); + await startTransport(transport); const enabling = transport.setVideoEnabled(true); await vi.waitFor(() => expect(getUserMedia).toHaveBeenCalledTimes(2)); transport.stop(); @@ -316,6 +339,45 @@ describe("Google Live Video Talk", () => { expect(FakeGoogleLiveWebSocket.instance?.readyState).toBe(3); }); + it("finishes active camera cleanup when the stream callback throws", async () => { + const audioStop = vi.fn(); + const videoStop = vi.fn(); + const audioTrack = { stop: audioStop } as unknown as MediaStreamTrack; + const videoTrack = Object.assign(new EventTarget(), { + stop: videoStop, + readyState: "live", + enabled: true, + muted: false, + }) as unknown as MediaStreamTrack; + const audio = { + getAudioTracks: () => [audioTrack], + getTracks: () => [audioTrack], + } as unknown as MediaStream; + const camera = { + getVideoTracks: () => [videoTrack], + getTracks: () => [videoTrack], + } as unknown as MediaStream; + vi.stubGlobal("navigator", { + mediaDevices: { + getUserMedia: vi.fn().mockResolvedValueOnce(audio).mockResolvedValueOnce(camera), + }, + }); + vi.spyOn(HTMLMediaElement.prototype, "play").mockResolvedValue(undefined); + const onVideoStream = vi.fn((stream: MediaStream | null) => { + if (!stream) { + throw new Error("stream callback failed"); + } + }); + const transport = createTransport({ onVideoStream }); + const ws = await startTransport(transport); + await transport.setVideoEnabled(true); + + expect(() => transport.stop()).toThrow("stream callback failed"); + expect(audioStop).toHaveBeenCalledOnce(); + expect(videoStop).toHaveBeenCalledOnce(); + expect(ws.readyState).toBe(3); + }); + it("switches an active camera and keeps video frame capture running", async () => { const audioTrack = { stop: vi.fn() } as unknown as MediaStreamTrack; const frontStop = vi.fn(); @@ -355,7 +417,7 @@ describe("Google Live Video Talk", () => { const onVideoStream = vi.fn(); const transport = createTransport({ onVideoStream }, "front"); - await transport.start(); + await startTransport(transport); await transport.setVideoEnabled(true); await transport.switchCamera("back"); @@ -403,24 +465,18 @@ describe("Google Live Video Talk", () => { const onVideoStream = vi.fn(); const transport = createTransport({ onStatus, onVideoStream }); - await transport.start(); + const { start, ws } = await beginTransport(transport); await transport.setVideoEnabled(true); - const ws = FakeGoogleLiveWebSocket.instance; - if (!ws) { - throw new Error("missing Google Live WebSocket"); - } ws.emitOpen(); + const rejected = expect(start).rejects.toThrow("Realtime connection timed out after 30000ms"); await vi.advanceTimersByTimeAsync(30_000); - expect(onStatus).toHaveBeenCalledExactlyOnceWith( - "error", - "Realtime connection timed out after 30000ms", - ); + await rejected; + expect(onStatus).not.toHaveBeenCalled(); expect(audioStop).toHaveBeenCalledOnce(); expect(videoStop).toHaveBeenCalledOnce(); expect(getUserMedia).toHaveBeenNthCalledWith(2, { video: true }); - expect(onVideoStream).toHaveBeenNthCalledWith(1, camera); - expect(onVideoStream).toHaveBeenLastCalledWith(null); + expect(onVideoStream).not.toHaveBeenCalled(); expect(ws.readyState).toBe(3); }); }); From 63cb1e9d2e94e98402acfdc2bb01ccc009cab7ef Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 17:29:15 +0800 Subject: [PATCH 7/9] fix(talk): handle reentrant Google activation cancellation --- .../realtime-talk-google-live-lifecycle.ts | 2 +- .../realtime-talk-google-live-timeout.test.ts | 26 +++++++++++++++ .../pages/chat/realtime-talk-google-live.ts | 32 +++++++++++-------- 3 files changed, 45 insertions(+), 15 deletions(-) diff --git a/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts b/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts index 1d12604af401..00679e274405 100644 --- a/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts +++ b/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts @@ -32,7 +32,7 @@ export function buildGoogleLiveUrl(session: RealtimeTalkJsonPcmWebSocketSessionR return url.toString(); } -export type GoogleLiveConnectionState = +type GoogleLiveConnectionState = | "idle" | "connecting" | "ready" diff --git a/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts b/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts index e3d1d62d4b5a..dfe0c20d66ae 100644 --- a/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts +++ b/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts @@ -67,6 +67,15 @@ class FakeAudioContext { createGain() { return { connect() {}, disconnect() {}, gain: { value: 1 } }; } + + createAnalyser() { + return { + fftSize: 0, + smoothingTimeConstant: 0, + disconnect() {}, + getFloatTimeDomainData: (samples: Float32Array) => samples.fill(0.25), + }; + } } function createSession(): RealtimeTalkJsonPcmWebSocketSessionResult { @@ -247,6 +256,23 @@ describe("Google Live setup timeout", () => { expect(audioContexts.every((context) => context.close.mock.calls.length === 1)).toBe(true); }); + it("reclaims the meter when an input-level callback cancels activation", async () => { + let stopDuringActivation = () => undefined; + const onInputLevel = vi.fn(() => stopDuringActivation()); + const transport = createTransport({ onInputLevel }); + stopDuringActivation = () => transport.stop({ emitClosed: false }); + const { start, socket } = await beginTransport(transport); + socket.emitOpen(); + socket.emitMessage({ setupComplete: {} }); + await expect(start).resolves.toBe("ready"); + + expect(() => transport.activate()).toThrow("Google Live transport activation cancelled"); + expect(vi.getTimerCount()).toBe(0); + expect(stopInputTrack).toHaveBeenCalledOnce(); + expect(socket.readyState).toBe(3); + expect(audioContexts.every((context) => context.close.mock.calls.length === 1)).toBe(true); + }); + it("clears the deadline after Google setup completes", async () => { const onStatus = vi.fn(); const transport = createTransport({ onStatus }); diff --git a/ui/src/pages/chat/realtime-talk-google-live.ts b/ui/src/pages/chat/realtime-talk-google-live.ts index 38160a287524..9f9f73b31849 100644 --- a/ui/src/pages/chat/realtime-talk-google-live.ts +++ b/ui/src/pages/chat/realtime-talk-google-live.ts @@ -209,28 +209,26 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { } try { this.ctx.callbacks.onStatus?.("listening"); - if (this.closed) { - return; - } + this.assertActivationCurrent(); this.emitTalkEvent({ type: "session.ready" }); - if (this.closed) { - return; - } + this.assertActivationCurrent(); if (this.ctx.callbacks.onInputLevel && this.media && this.inputContext) { - this.inputMeter = new RealtimeTalkMediaStreamMeter(this.ctx.callbacks.onInputLevel); - this.inputMeter.start(this.media, this.inputContext); - } - if (this.closed) { - return; + const inputMeter = new RealtimeTalkMediaStreamMeter(this.ctx.callbacks.onInputLevel); + this.inputMeter = inputMeter; + inputMeter.start(this.media, this.inputContext); + if (this.closed || !this.lifecycle.isActive || this.inputMeter !== inputMeter) { + // start() publishes synchronously before installing its interval. A + // reentrant stop must reclaim the interval that start() installs next. + inputMeter.stop(false); + } + this.assertActivationCurrent(); } this.startMicrophonePump(); if (this.camera.stream && !this.cameraPublished) { this.cameraPublished = true; this.ctx.callbacks.onVideoStream?.(this.camera.stream); } - if (this.closed) { - return; - } + this.assertActivationCurrent(); this.startVideoFrames(); } catch (error) { try { @@ -242,6 +240,12 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { } } + private assertActivationCurrent(): void { + if (this.closed || !this.lifecycle.isActive) { + throw new Error("Google Live transport activation cancelled"); + } + } + async setVideoEnabled(enabled: boolean): Promise { await this.camera.setEnabled(enabled); } From 9c6132ef8c0d0971df2e4c4bec2433cd1461571a Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 17:34:50 +0800 Subject: [PATCH 8/9] test(talk): type activation cancellation callback --- ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts b/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts index dfe0c20d66ae..1ed89bd74d44 100644 --- a/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts +++ b/ui/src/pages/chat/realtime-talk-google-live-timeout.test.ts @@ -257,7 +257,7 @@ describe("Google Live setup timeout", () => { }); it("reclaims the meter when an input-level callback cancels activation", async () => { - let stopDuringActivation = () => undefined; + let stopDuringActivation: () => void = () => undefined; const onInputLevel = vi.fn(() => stopDuringActivation()); const transport = createTransport({ onInputLevel }); stopDuringActivation = () => transport.stop({ emitClosed: false }); From b993cf1b99659bfc928c8057c2b12bc14f278c4c Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 17:42:18 +0800 Subject: [PATCH 9/9] refactor(talk): centralize Google Live cleanup errors --- .../realtime-talk-google-live-lifecycle.ts | 7 ++++-- .../pages/chat/realtime-talk-google-live.ts | 24 +++++++------------ 2 files changed, 13 insertions(+), 18 deletions(-) diff --git a/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts b/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts index 00679e274405..6a4cd8fa99e0 100644 --- a/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts +++ b/ui/src/pages/chat/realtime-talk-google-live-lifecycle.ts @@ -131,12 +131,15 @@ export class GoogleLiveConnectionLifecycle { } export function runRealtimeTalkCleanup(steps: Array<() => void>): void { - let firstError: unknown; + let firstError: Error | undefined; for (const step of steps) { try { step(); } catch (error) { - firstError ??= error; + firstError ??= + error instanceof Error + ? error + : new Error("Realtime Talk cleanup failed", { cause: error }); } } if (firstError) { diff --git a/ui/src/pages/chat/realtime-talk-google-live.ts b/ui/src/pages/chat/realtime-talk-google-live.ts index 9f9f73b31849..d507c45ec251 100644 --- a/ui/src/pages/chat/realtime-talk-google-live.ts +++ b/ui/src/pages/chat/realtime-talk-google-live.ts @@ -258,22 +258,14 @@ export class GoogleLiveRealtimeTalkTransport implements RealtimeTalkTransport { const emitClosed = !this.closed && this.lifecycle.isActive && options?.emitClosed !== false; this.closed = true; this.lifecycle.cancel(); - let firstError: unknown; - try { - if (emitClosed) { - this.emitTalkEvent({ type: "session.closed", final: true }); - } - } catch (error) { - firstError = error; - } - try { - this.releaseResources(); - } catch (error) { - firstError ??= error; - } - if (firstError) { - throw firstError; - } + runRealtimeTalkCleanup([ + () => { + if (emitClosed) { + this.emitTalkEvent({ type: "session.closed", final: true }); + } + }, + () => this.releaseResources(), + ]); } private releaseResources(): void {