From 38448750ebf55b5964393dea61c19fe9fdd96622 Mon Sep 17 00:00:00 2001 From: Wynne668 Date: Thu, 16 Jul 2026 23:27:36 +0800 Subject: [PATCH] fix(openai): stop realtime voice reconnects promptly on close (#108209) * fix(openai): cancel realtime reconnect backoff on close * test(openai): cover reconnect after bridge close --------- Co-authored-by: Peter Steinberger --- .../openai/realtime-voice-provider.test.ts | 49 +++++++++++++++++++ extensions/openai/realtime-voice-provider.ts | 21 ++++++-- 2 files changed, 66 insertions(+), 4 deletions(-) diff --git a/extensions/openai/realtime-voice-provider.test.ts b/extensions/openai/realtime-voice-provider.test.ts index 04ed436e4802..00d700e97715 100644 --- a/extensions/openai/realtime-voice-provider.test.ts +++ b/extensions/openai/realtime-voice-provider.test.ts @@ -1101,6 +1101,55 @@ describe("buildOpenAIRealtimeVoiceProvider", () => { bridge.close(); }); + it("cancels a pending reconnect and allows a later explicit connect", async () => { + vi.useFakeTimers(); + const provider = buildOpenAIRealtimeVoiceProvider(); + const onError = vi.fn(); + const bridge = provider.createBridge({ + providerConfig: { apiKey: "sk-test" }, // pragma: allowlist secret + onAudio: vi.fn(), + onClearAudio: vi.fn(), + onError, + }); + const connecting = bridge.connect(); + const socket = FakeWebSocket.instances[0]; + if (!socket) { + throw new Error("expected bridge to create a websocket"); + } + + socket.readyState = FakeWebSocket.OPEN; + socket.emit("open"); + socket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + await connecting; + + socket.readyState = FakeWebSocket.CLOSED; + socket.emit("close", 1006, Buffer.from("transient drop")); + await vi.advanceTimersByTimeAsync(0); + expect(vi.getTimerCount()).toBe(1); + + bridge.close(); + await vi.advanceTimersByTimeAsync(0); + + expect(vi.getTimerCount()).toBe(0); + expect(FakeWebSocket.instances).toHaveLength(1); + expect(onError).not.toHaveBeenCalled(); + + const reconnecting = bridge.connect(); + const reconnectedSocket = FakeWebSocket.instances[1]; + if (!reconnectedSocket) { + throw new Error("expected bridge to reconnect after close"); + } + reconnectedSocket.readyState = FakeWebSocket.OPEN; + reconnectedSocket.emit("open"); + reconnectedSocket.emit("message", Buffer.from(JSON.stringify({ type: "session.updated" }))); + await reconnecting; + + expect(bridge.isConnected()).toBe(true); + expect(FakeWebSocket.instances).toHaveLength(2); + expect(onError).not.toHaveBeenCalled(); + bridge.close(); + }); + it("keeps Azure deployment bridges on deployment-compatible session payloads", async () => { const provider = buildOpenAIRealtimeVoiceProvider(); const bridge = provider.createBridge({ diff --git a/extensions/openai/realtime-voice-provider.ts b/extensions/openai/realtime-voice-provider.ts index 9068019f154d..9943c580b21f 100644 --- a/extensions/openai/realtime-voice-provider.ts +++ b/extensions/openai/realtime-voice-provider.ts @@ -27,7 +27,7 @@ import { REALTIME_VOICE_AUDIO_FORMAT_G711_ULAW_8KHZ, REALTIME_VOICE_AUDIO_FORMAT_PCM16_24KHZ, } from "openclaw/plugin-sdk/realtime-voice"; -import { warn } from "openclaw/plugin-sdk/runtime-env"; +import { sleepWithAbort, warn } from "openclaw/plugin-sdk/runtime-env"; import { normalizeResolvedSecretInputString, normalizeSecretInputString, @@ -506,6 +506,7 @@ class OpenAIRealtimeVoiceBridge implements RealtimeVoiceBridge { private sessionReadyFired = false; private reconnectReason: string | undefined; private activeConnectionReason: string | undefined; + private reconnectAbortController = new AbortController(); private readonly audioFormat: RealtimeVoiceAudioFormat; constructor(private readonly config: OpenAIRealtimeVoiceBridgeConfig) { @@ -514,6 +515,9 @@ class OpenAIRealtimeVoiceBridge implements RealtimeVoiceBridge { async connect(): Promise { this.intentionallyClosed = false; + if (this.reconnectAbortController.signal.aborted) { + this.reconnectAbortController = new AbortController(); + } this.reconnectAttempts = 0; await this.doConnect(); } @@ -587,6 +591,9 @@ class OpenAIRealtimeVoiceBridge implements RealtimeVoiceBridge { close(): void { this.intentionallyClosed = true; + // The bridge owns both its active socket and reconnect delay; canceling + // both keeps terminal close from retaining callbacks for the full backoff. + this.reconnectAbortController.abort(); this.connected = false; this.sessionConfigured = false; if (this.ws) { @@ -881,9 +888,15 @@ class OpenAIRealtimeVoiceBridge implements RealtimeVoiceBridge { type: "session.reconnect.scheduled", detail: `reason=${reason} attempt=${attempt} delayMs=${delay}`, }); - await new Promise((resolve) => { - setTimeout(resolve, delay); - }); + const reconnectSignal = this.reconnectAbortController.signal; + try { + await sleepWithAbort(delay, reconnectSignal); + } catch (error) { + if (!reconnectSignal.aborted) { + throw error; + } + return; + } if (this.intentionallyClosed) { return; }