From a38dc473dc7a46bdf9062cc17633c3385758e89c Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sun, 2 Aug 2026 01:06:14 +0800 Subject: [PATCH] fix(google): make continuity reset generation-idempotent --- .../google/realtime-voice-provider.test.ts | 83 +++++++++++++++++++ extensions/google/realtime-voice-provider.ts | 9 +- 2 files changed, 91 insertions(+), 1 deletion(-) diff --git a/extensions/google/realtime-voice-provider.test.ts b/extensions/google/realtime-voice-provider.test.ts index 8cc472ca6123..c01314c84ccb 100644 --- a/extensions/google/realtime-voice-provider.test.ts +++ b/extensions/google/realtime-voice-provider.test.ts @@ -1071,6 +1071,89 @@ describe("buildGoogleRealtimeVoiceProvider", () => { ); }); + it("emits one continuity reset across failed fresh reconnect attempts", async () => { + vi.useFakeTimers(); + const provider = buildGoogleRealtimeVoiceProvider(); + const onClose = vi.fn(); + const onEvent = vi.fn(); + const bridge = provider.createBridge({ + providerConfig: { apiKey: "gemini-key", sessionResumption: false }, + onAudio: vi.fn(), + onClearAudio: vi.fn(), + onClose, + onEvent, + }); + + await bridge.connect(); + const firstSession = lastConnectParams().callbacks; + firstSession.onmessage({ setupComplete: {} }); + connectMock + .mockRejectedValueOnce(new Error("connect failed 1")) + .mockRejectedValueOnce(new Error("connect failed 2")) + .mockRejectedValueOnce(new Error("connect failed 3")); + firstSession.onclose({ code: 1011, reason: "temporary" }); + + await vi.advanceTimersByTimeAsync(1_750); + + expect(connectMock).toHaveBeenCalledTimes(4); + expect(onEvent.mock.calls).toEqual([ + [{ direction: "client", type: "session.continuity.reset" }], + ]); + expect(onClose).toHaveBeenCalledWith("error"); + }); + + it("rearms continuity reset after pre-return setup selects a fresh session", async () => { + vi.useFakeTimers(); + const pendingSession = createDeferred(); + const freshSession = createMockGoogleLiveSession(); + connectMock + .mockReturnValueOnce(Promise.resolve(session)) + .mockReturnValueOnce(pendingSession.promise); + const provider = buildGoogleRealtimeVoiceProvider(); + const onEvent = vi.fn(); + const onReady = vi.fn(); + const onTranscript = vi.fn(); + const bridge = provider.createBridge({ + providerConfig: { apiKey: "gemini-key", sessionResumption: false }, + onAudio: vi.fn(), + onClearAudio: vi.fn(), + onEvent, + onReady, + onTranscript, + }); + + await bridge.connect(); + const firstCallbacks = lastConnectParams().callbacks; + firstCallbacks.onopen(); + firstCallbacks.onmessage({ setupComplete: {} }); + firstCallbacks.onclose({ code: 1011, reason: "temporary" }); + await vi.advanceTimersByTimeAsync(250); + + const freshCallbacks = lastConnectParams().callbacks; + expect(onEvent).toHaveBeenCalledTimes(1); + freshCallbacks.onopen(); + freshCallbacks.onmessage({ + setupComplete: {}, + serverContent: { inputTranscription: { text: "Fresh partial " } }, + }); + expect(onReady).toHaveBeenCalledTimes(1); + + pendingSession.resolve(freshSession); + await vi.waitFor(() => { + expect(onReady).toHaveBeenCalledTimes(2); + }); + freshCallbacks.onclose({ code: 1011, reason: "temporary again" }); + await vi.advanceTimersByTimeAsync(250); + + expect(onEvent).toHaveBeenCalledTimes(2); + lastConnectParams().callbacks.onmessage({ + serverContent: { inputTranscription: { text: "Next", finished: true } }, + }); + expect(onTranscript.mock.calls.filter((call) => call[2] === true)).toEqual([ + ["user", "Next", true], + ]); + }); + it("waits for the returned session after setup completion before activating", async () => { const pendingSession = createDeferred(); const connectedSession = createMockGoogleLiveSession(); diff --git a/extensions/google/realtime-voice-provider.ts b/extensions/google/realtime-voice-provider.ts index dbfeec2c4878..89c1715d33f4 100644 --- a/extensions/google/realtime-voice-provider.ts +++ b/extensions/google/realtime-voice-provider.ts @@ -482,6 +482,7 @@ class GoogleRealtimeVoiceBridge implements RealtimeVoiceBridge { private reconnectAttempts = 0; private reconnectTimer: ReturnType | undefined; private hasConnectedSession = false; + private continuityResetEmitted = false; private terminalError: Error | undefined; private closeNotified = false; private connectionOwner: GoogleLiveConnectionAttempt | undefined; @@ -532,10 +533,11 @@ class GoogleRealtimeVoiceBridge implements RealtimeVoiceBridge { private async connectOwned(attempt: GoogleLiveConnectionAttempt): Promise { const canResumeSession = this.config.sessionResumption !== false && Boolean(this.resumptionHandle); - if (this.hasConnectedSession && !canResumeSession) { + if (this.hasConnectedSession && !canResumeSession && !this.continuityResetEmitted) { // An unfinished recognition hypothesis cannot cross into a fresh server session. // Notify consumers before connect because the SDK can replay fresh-session // callbacks before its connect promise returns. + this.continuityResetEmitted = true; this.resetPendingTranscripts(); this.config.onEvent?.({ direction: "client", @@ -868,6 +870,11 @@ class GoogleRealtimeVoiceBridge implements RealtimeVoiceBridge { } private handleSetupComplete(): void { + if (!this.setupCompleteReceived) { + // setupComplete proves Google selected a new server session. A later + // continuity loss therefore owns a new reset generation. + this.continuityResetEmitted = false; + } this.setupCompleteReceived = true; this.maybeActivateSession(); }