diff --git a/src/talk/session-runtime.test.ts b/src/talk/session-runtime.test.ts index 8d944bd92ba8..bf5f58f878cd 100644 --- a/src/talk/session-runtime.test.ts +++ b/src/talk/session-runtime.test.ts @@ -343,6 +343,167 @@ describe("realtime voice bridge session runtime", () => { expect(onError).not.toHaveBeenCalled(); }); + it("permanently closes once while preserving synchronous transcript flush", async () => { + let callbacks: Parameters[0] | undefined; + const close = vi.fn(() => { + callbacks?.onTranscript?.("assistant", "final transcript", true); + }); + const connect = vi.fn(async () => {}); + const sendProviderAudio = vi.fn(); + const sendSinkAudio = vi.fn(); + const onTranscript = vi.fn(); + const provider: RealtimeVoiceProviderPlugin = { + id: "test", + label: "Test", + isConfigured: () => true, + createBridge: (request) => { + callbacks = request; + return makeBridge({ close, connect, sendAudio: sendProviderAudio }); + }, + }; + const session = createRealtimeVoiceBridgeSession({ + provider, + providerConfig: {}, + audioSink: { sendAudio: sendSinkAudio }, + onTranscript, + }); + + session.close(); + session.close(); + session.sendAudio(Buffer.from("late-input")); + callbacks?.onAudio(Buffer.from("late-output")); + await expect(session.connect()).rejects.toThrow("Realtime voice session is closed"); + + expect(close).toHaveBeenCalledTimes(1); + expect(connect).not.toHaveBeenCalled(); + expect(sendProviderAudio).not.toHaveBeenCalled(); + expect(sendSinkAudio).not.toHaveBeenCalled(); + expect(onTranscript).toHaveBeenCalledExactlyOnceWith("assistant", "final transcript", true); + }); + + it("stops audio admission after provider close and still closes the provider once", () => { + let callbacks: Parameters[0] | undefined; + const close = vi.fn(); + const sendProviderAudio = vi.fn(); + const sendSinkAudio = vi.fn(); + const onClose = vi.fn(); + const provider: RealtimeVoiceProviderPlugin = { + id: "test", + label: "Test", + isConfigured: () => true, + createBridge: (request) => { + callbacks = request; + return makeBridge({ close, sendAudio: sendProviderAudio }); + }, + }; + const session = createRealtimeVoiceBridgeSession({ + provider, + providerConfig: {}, + audioSink: { sendAudio: sendSinkAudio }, + onClose, + }); + + callbacks?.onClose?.("completed"); + callbacks?.onClose?.("completed"); + session.sendAudio(Buffer.from("late-input")); + callbacks?.onAudio(Buffer.from("late-output")); + session.close(); + session.close(); + + expect(onClose).toHaveBeenCalledTimes(1); + expect(close).toHaveBeenCalledTimes(1); + expect(sendProviderAudio).not.toHaveBeenCalled(); + expect(sendSinkAudio).not.toHaveBeenCalled(); + }); + + it("reopens audio and close reporting for an explicit connection generation", async () => { + let callbacks: Parameters[0] | undefined; + const connect = vi.fn(async () => {}); + const sendProviderAudio = vi.fn(); + const sendSinkAudio = vi.fn(); + const onClose = vi.fn(); + const onReady = vi.fn(); + const provider: RealtimeVoiceProviderPlugin = { + id: "test", + label: "Test", + isConfigured: () => true, + createBridge: (request) => { + callbacks = request; + request.onClose?.("error"); + return makeBridge({ connect, sendAudio: sendProviderAudio }); + }, + }; + const session = createRealtimeVoiceBridgeSession({ + provider, + providerConfig: {}, + audioSink: { sendAudio: sendSinkAudio }, + onClose, + onReady, + }); + + session.sendAudio(Buffer.from("closed-input")); + callbacks?.onAudio(Buffer.from("closed-output")); + callbacks?.onClose?.("error"); + callbacks?.onReady?.(); + expect(onReady).not.toHaveBeenCalled(); + + await session.connect(); + callbacks?.onReady?.(); + session.sendAudio(Buffer.from("next-input")); + callbacks?.onAudio(Buffer.from("next-output")); + callbacks?.onClose?.("completed"); + callbacks?.onClose?.("completed"); + await expect(session.connect()).rejects.toThrow("Realtime voice connection is closed"); + + expect(connect).toHaveBeenCalledTimes(1); + expect(onReady).toHaveBeenCalledWith(session); + expect(onClose).toHaveBeenNthCalledWith(1, "error"); + expect(onClose).toHaveBeenNthCalledWith(2, "completed"); + expect(onClose).toHaveBeenCalledTimes(2); + expect(sendProviderAudio).toHaveBeenCalledExactlyOnceWith(Buffer.from("next-input")); + expect(sendSinkAudio).toHaveBeenCalledExactlyOnceWith(Buffer.from("next-output")); + }); + + it("rejects reconnect and ignores tool failures after an established provider close", async () => { + let callbacks: Parameters[0] | undefined; + let rejectToolCall: ((error: Error) => void) | undefined; + const connect = vi.fn(async () => {}); + const onError = vi.fn(); + const provider: RealtimeVoiceProviderPlugin = { + id: "test", + label: "Test", + isConfigured: () => true, + createBridge: (request) => { + callbacks = request; + return makeBridge({ connect }); + }, + }; + const session = createRealtimeVoiceBridgeSession({ + provider, + providerConfig: {}, + audioSink: { sendAudio: vi.fn() }, + onToolCall: () => + new Promise((_resolve, reject) => { + rejectToolCall = reject; + }), + onError, + }); + + callbacks?.onToolCall?.({ + itemId: "item-1", + callId: "call-1", + name: "lookup", + args: {}, + }); + callbacks?.onClose?.("error"); + await expect(session.connect()).rejects.toThrow("Realtime voice connection is closed"); + rejectToolCall?.(new Error("late tool callback failure")); + await Promise.resolve(); + + expect(connect).not.toHaveBeenCalled(); + expect(onError).not.toHaveBeenCalled(); + }); + it("forwards tool result continuation options and async acceptance to the provider bridge", () => { const acceptance = Promise.resolve(); const submitToolResult = vi.fn(() => acceptance); diff --git a/src/talk/session-runtime.ts b/src/talk/session-runtime.ts index b46c2aace94a..68b9a5c28593 100644 --- a/src/talk/session-runtime.ts +++ b/src/talk/session-runtime.ts @@ -78,6 +78,8 @@ export type RealtimeVoiceBridgeSessionParams = { onClose?: (reason: RealtimeVoiceCloseReason) => void; }; +type RealtimeVoiceSessionPhase = "admitting" | "provider-terminal" | "disposed"; + /** * Creates a realtime voice bridge session and wires provider events to the configured audio sink. */ @@ -85,7 +87,12 @@ export function createRealtimeVoiceBridgeSession( params: RealtimeVoiceBridgeSessionParams, ): RealtimeVoiceBridgeSession { const bridgeRef: { current?: RealtimeVoiceBridge } = {}; - let isActive = true; + // Local disposal owns provider cleanup. Only a terminal callback fired before bridge + // adoption may reopen; adopted bridges own reconnects and stale-event fencing internally. + let phase: RealtimeVoiceSessionPhase = "admitting"; + let terminalBeforeBridgeAdoption = false; + let closeReported = false; + const isAdmitting = () => phase === "admitting"; const requireBridge = () => { if (!bridgeRef.current) { throw new Error("Realtime voice bridge is not ready"); @@ -100,12 +107,32 @@ export function createRealtimeVoiceBridgeSession( }, acknowledgeMark: (markName) => requireBridge().acknowledgeMark(markName), close: () => { + if (phase === "disposed") { + return; + } const bridge = requireBridge(); - isActive = false; + phase = "disposed"; bridge.close(); }, - connect: () => requireBridge().connect(), - sendAudio: (audio) => requireBridge().sendAudio(audio), + connect: () => { + if (phase === "disposed") { + return Promise.reject(new Error("Realtime voice session is closed")); + } + if (phase === "provider-terminal") { + if (!terminalBeforeBridgeAdoption) { + return Promise.reject(new Error("Realtime voice connection is closed")); + } + terminalBeforeBridgeAdoption = false; + phase = "admitting"; + closeReported = false; + } + return requireBridge().connect(); + }, + sendAudio: (audio) => { + if (isAdmitting()) { + requireBridge().sendAudio(audio); + } + }, sendUserMessage: (text) => requireBridge().sendUserMessage?.(text), handleBargeIn: (options) => requireBridge().handleBargeIn?.(options), setMediaTimestamp: (ts) => requireBridge().setMediaTimestamp(ts), @@ -118,11 +145,13 @@ export function createRealtimeVoiceBridgeSession( }, triggerGreeting: (instructions) => requireBridge().triggerGreeting?.(instructions), }; - const canSendAudio = () => params.audioSink.isOpen?.() ?? true; + // Session inactivity is the shared admission boundary for both audio directions. + // Provider and transport callbacks may still race after close, but cannot retain new audio. + const canSendAudio = () => isAdmitting() && (params.audioSink.isOpen?.() ?? true); const reportCallbackError = (error: unknown) => { // Async tool handlers can settle after the provider closes. Once inactive, no // callback may report stale failures into the next session lifecycle. - if (!isActive) { + if (!isAdmitting()) { return; } try { @@ -167,7 +196,7 @@ export function createRealtimeVoiceBridgeSession( onTranscript: params.onTranscript, onEvent: params.onEvent, onToolCall: (event) => { - if (!bridgeRef.current || !isActive) { + if (!bridgeRef.current || !isAdmitting()) { return; } try { @@ -180,7 +209,7 @@ export function createRealtimeVoiceBridgeSession( } }, onReady: () => { - if (!bridgeRef.current) { + if (!bridgeRef.current || !isAdmitting()) { return; } if (params.triggerGreetingOnReady) { @@ -190,7 +219,16 @@ export function createRealtimeVoiceBridgeSession( }, onError: params.onError, onClose: (reason) => { - isActive = false; + if (!bridgeRef.current) { + terminalBeforeBridgeAdoption = true; + } + if (phase !== "disposed") { + phase = "provider-terminal"; + } + if (closeReported) { + return; + } + closeReported = true; params.onClose?.(reason); }, });