From c6b3680ecf2350f6e92133841729131e0d8854d0 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Fri, 31 Jul 2026 13:46:27 -0700 Subject: [PATCH 1/5] fix(openai): preserve gateway microphone startup audio --- ...ealtime-quicksilver-gateway-bridge.test.ts | 145 ++++++++++++++++++ .../realtime-quicksilver-gateway-bridge.ts | 22 ++- 2 files changed, 166 insertions(+), 1 deletion(-) diff --git a/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts b/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts index c35b58dd6392..0c5e54d4d670 100644 --- a/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts +++ b/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts @@ -522,6 +522,151 @@ describe("GPT-Live werift audio peer", () => { }); describe("GPT-Live gateway relay bridge", () => { + function createPendingPeerBridge() { + let resolvePeer: ((peer: OpenAIQuicksilverAudioPeerContract) => void) | undefined; + let rejectPeer: ((error: Error) => void) | undefined; + const peerPromise = new Promise((resolve, reject) => { + resolvePeer = resolve; + rejectPeer = reject; + }); + const peer = { + createOffer: vi.fn(async () => "v=offer\r\n"), + applyAnswer: vi.fn(async () => undefined), + sendAudio: vi.fn(), + close: vi.fn(), + } satisfies OpenAIQuicksilverAudioPeerContract; + const bridge = new OpenAIQuicksilverGatewayBridge({ + providerConfig: {}, + model: "gpt-live-1-codex", + voice: "marin", + audioFormat: { encoding: "pcm16", sampleRateHz: 24_000, channels: 1 }, + onAudio: vi.fn(), + onClearAudio: vi.fn(), + runAgentConsult: vi.fn(async () => ({ text: "done" })), + logger: { debug: vi.fn(), warn: vi.fn() }, + resolveAuth: vi.fn(async () => ({ + type: "oauth" as const, + token: "oauth-token", + accountId: "account-1", + })), + createPeer: vi.fn(() => peerPromise), + fetchImpl: vi.fn(async () => createCallResponse("v=answer\r\n", "rtc_pending_audio")), + webSocketFactory: () => new FakeSocket(), + }); + const connection = bridge.connect(); + return { + bridge, + connection, + peer, + rejectPeer: (error: Error) => rejectPeer?.(error), + resolvePeer: () => resolvePeer?.(peer), + }; + } + + it("preserves caller-owned microphone frames while the media peer is starting", async () => { + const { bridge, connection, peer, resolvePeer } = createPendingPeerBridge(); + try { + const source = Buffer.from([0x7f, 0x41]); + bridge.sendAudio(source); + source.fill(0); + bridge.sendAudio(Buffer.from([0x22, 0x23])); + + resolvePeer(); + await connection; + + expect(peer.sendAudio.mock.calls.map(([audio]) => audio)).toEqual([ + Buffer.from([0x7f, 0x41]), + Buffer.from([0x22, 0x23]), + ]); + bridge.sendAudio(Buffer.from([0x30, 0x31])); + expect(peer.sendAudio).toHaveBeenCalledTimes(3); + } finally { + bridge.close(); + } + }); + + it("bounds queued microphone bytes before copying rejected audio", async () => { + const { bridge, connection, peer, resolvePeer } = createPendingPeerBridge(); + try { + bridge.sendAudio(Buffer.alloc(512 * 1024, 0x01)); + bridge.sendAudio(Buffer.alloc(512 * 1024, 0x02)); + const overflow = Buffer.alloc(1, 0x03); + const copy = vi.spyOn(Buffer, "from"); + try { + bridge.sendAudio(overflow); + expect(copy).not.toHaveBeenCalled(); + } finally { + copy.mockRestore(); + } + + resolvePeer(); + await connection; + + expect(peer.sendAudio).toHaveBeenCalledTimes(2); + expect(peer.sendAudio.mock.calls.map(([audio]) => audio.byteLength)).toEqual([ + 512 * 1024, + 512 * 1024, + ]); + } finally { + bridge.close(); + } + }); + + it("bounds queued microphone frame count before copying rejected audio", async () => { + const { bridge, connection, peer, resolvePeer } = createPendingPeerBridge(); + try { + for (let index = 0; index < 320; index += 1) { + bridge.sendAudio(Buffer.alloc(2, index)); + } + const overflow = Buffer.alloc(2, 0xff); + const copy = vi.spyOn(Buffer, "from"); + try { + bridge.sendAudio(overflow); + expect(copy).not.toHaveBeenCalled(); + } finally { + copy.mockRestore(); + } + + resolvePeer(); + await connection; + + expect(peer.sendAudio).toHaveBeenCalledTimes(320); + expect(peer.sendAudio.mock.calls.at(-1)?.[0]).toEqual(Buffer.alloc(2, 319)); + } finally { + bridge.close(); + } + }); + + it("discards queued microphone audio when closed before the media peer resolves", async () => { + const { bridge, connection, peer, resolvePeer } = createPendingPeerBridge(); + bridge.sendAudio(Buffer.from([0x41, 0x42])); + bridge.close(); + resolvePeer(); + + await expect(connection).rejects.toThrow("GPT-Live gateway relay bridge closed"); + await vi.waitFor(() => expect(peer.close).toHaveBeenCalledOnce()); + expect(peer.sendAudio).not.toHaveBeenCalled(); + bridge.sendAudio(Buffer.from([0x43, 0x44])); + expect(peer.sendAudio).not.toHaveBeenCalled(); + }); + + it("discards queued microphone audio when media peer creation fails", async () => { + const { bridge, connection, peer, rejectPeer } = createPendingPeerBridge(); + const pendingAudioState = bridge as unknown as { + pendingAudio: Buffer[]; + pendingAudioBytes: number; + }; + bridge.sendAudio(Buffer.from([0x41, 0x42])); + rejectPeer(new Error("media peer unavailable")); + + await expect(connection).rejects.toThrow("media peer unavailable"); + expect(pendingAudioState.pendingAudio).toEqual([]); + expect(pendingAudioState.pendingAudioBytes).toBe(0); + bridge.sendAudio(Buffer.from([0x43, 0x44])); + expect(pendingAudioState.pendingAudio).toEqual([]); + expect(peer.sendAudio).not.toHaveBeenCalled(); + }); + it("closes a sideband that opens in the abort handoff", async () => { const controller = new AbortController(); const socket = new FakeSocket("manual"); diff --git a/extensions/openai/realtime-quicksilver-gateway-bridge.ts b/extensions/openai/realtime-quicksilver-gateway-bridge.ts index 73efd178c805..9e34ed4927b2 100644 --- a/extensions/openai/realtime-quicksilver-gateway-bridge.ts +++ b/extensions/openai/realtime-quicksilver-gateway-bridge.ts @@ -40,6 +40,7 @@ const RELAY_SAMPLE_RATE = 24_000; const QUICKSILVER_SESSION_TTL_MS = 30 * 60_000; const QUICKSILVER_CONNECT_TIMEOUT_MS = 30_000; const WEBSOCKET_OPEN = 1; +const MAX_PENDING_AUDIO = { chunks: 320, bytes: 1024 * 1024 }; function toError(error: unknown): Error { return error instanceof Error ? error : new Error(String(error)); @@ -164,6 +165,8 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { private closed = false; private closeNotified = false; private peer: OpenAIQuicksilverAudioPeerContract | undefined; + private pendingAudio: Buffer[] = []; + private pendingAudioBytes = 0; private ready = false; private sideband: ActiveSideband | undefined; private timer: ReturnType | undefined; @@ -181,7 +184,18 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { } sendAudio(audio: Buffer): void { - this.peer?.sendAudio(audio); + if (this.peer) { + this.peer.sendAudio(audio); + } else if ( + !this.closed && + !this.abortController.signal.aborted && + this.pendingAudio.length < MAX_PENDING_AUDIO.chunks && + this.pendingAudioBytes + audio.byteLength <= MAX_PENDING_AUDIO.bytes + ) { + // Relay capture starts before asynchronous peer creation and may recycle its input buffers. + this.pendingAudio.push(Buffer.from(audio)); + this.pendingAudioBytes += audio.byteLength; + } } setMediaTimestamp(_ts: number): void {} @@ -255,6 +269,10 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { () => undefined, ); this.peer = await waitForConnectStep(peerPromise, connectSignal); + for (const audio of this.pendingAudio.splice(0)) { + this.peer.sendAudio(audio); + } + this.pendingAudioBytes = 0; const offerSdp = await waitForConnectStep(this.peer.createOffer(), connectSignal); const auth = await waitForConnectStep(this.config.resolveAuth(), connectSignal); const requestIds = { @@ -535,6 +553,8 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { private releaseResources(): void { releaseOpenAIQuicksilverSession(this); this.connected = false; + this.pendingAudio = []; + this.pendingAudioBytes = 0; this.abortController.abort(new Error("GPT-Live gateway relay bridge closed")); this.consultController?.abort(new Error("GPT-Live delegation stopped")); this.consultController = undefined; From 2dbdf4f4ba001fda47325f01217f9eab4df1c92d Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 13:21:25 +0800 Subject: [PATCH 2/5] fix(openai): unify gateway microphone buffering --- .../realtime-quicksilver-audio-buffer.test.ts | 37 ++++ .../realtime-quicksilver-audio-buffer.ts | 27 +++ ...ealtime-quicksilver-gateway-bridge.test.ts | 161 ++++++++++++------ .../realtime-quicksilver-gateway-bridge.ts | 24 +-- .../realtime-quicksilver-peer.runtime.ts | 24 +-- 5 files changed, 187 insertions(+), 86 deletions(-) create mode 100644 extensions/openai/realtime-quicksilver-audio-buffer.test.ts create mode 100644 extensions/openai/realtime-quicksilver-audio-buffer.ts diff --git a/extensions/openai/realtime-quicksilver-audio-buffer.test.ts b/extensions/openai/realtime-quicksilver-audio-buffer.test.ts new file mode 100644 index 000000000000..af92af2dc172 --- /dev/null +++ b/extensions/openai/realtime-quicksilver-audio-buffer.test.ts @@ -0,0 +1,37 @@ +import { describe, expect, it } from "vitest"; +import { + appendOpenAIQuicksilverPendingAudio, + OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES, +} from "./realtime-quicksilver-audio-buffer.js"; + +describe("GPT-Live pending microphone audio", () => { + it("copies caller-owned PCM16 and drops an incomplete sample", () => { + const source = Buffer.from([0x01, 0x02, 0x03]); + const pending = appendOpenAIQuicksilverPendingAudio(Buffer.alloc(0), source); + source.fill(0xff); + + expect(pending).toEqual(Buffer.from([0x01, 0x02])); + }); + + it("appends audio in capture order while it fits", () => { + const first = appendOpenAIQuicksilverPendingAudio(Buffer.alloc(0), Buffer.from([0x01, 0x02])); + const second = appendOpenAIQuicksilverPendingAudio(first, Buffer.from([0x03, 0x04])); + + expect(second).toEqual(Buffer.from([0x01, 0x02, 0x03, 0x04])); + }); + + it("retains the newest bounded tail across existing and oversized input", () => { + const existing = Buffer.alloc(OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES, 0x01); + const appended = appendOpenAIQuicksilverPendingAudio(existing, Buffer.from([0x02, 0x02])); + const oversized = Buffer.alloc(OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES + 2, 0x03); + + expect(appended).toHaveLength(OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES); + expect(appended.subarray(0, -2).every((byte) => byte === 0x01)).toBe(true); + expect(appended.subarray(-2)).toEqual(Buffer.from([0x02, 0x02])); + expect( + appendOpenAIQuicksilverPendingAudio(appended, oversized).equals( + Buffer.alloc(OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES, 0x03), + ), + ).toBe(true); + }); +}); diff --git a/extensions/openai/realtime-quicksilver-audio-buffer.ts b/extensions/openai/realtime-quicksilver-audio-buffer.ts new file mode 100644 index 000000000000..5616cb7635aa --- /dev/null +++ b/extensions/openai/realtime-quicksilver-audio-buffer.ts @@ -0,0 +1,27 @@ +const RELAY_FRAME_SAMPLES = 480; +const MAX_PENDING_RELAY_FRAMES = 250; + +export const OPENAI_QUICKSILVER_RELAY_FRAME_BYTES = RELAY_FRAME_SAMPLES * 2; +// One five-second tail spans peer startup and the connected media pump. +// Keeping the newest PCM bounds latency without changing policy at adoption. +export const OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES = + OPENAI_QUICKSILVER_RELAY_FRAME_BYTES * MAX_PENDING_RELAY_FRAMES; + +export function appendOpenAIQuicksilverPendingAudio(pending: Buffer, incoming: Buffer): Buffer { + const evenLength = incoming.length - (incoming.length % 2); + if (evenLength === 0) { + return pending; + } + const audio = incoming.subarray(0, evenLength); + if (audio.length >= OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES) { + return Buffer.from(audio.subarray(audio.length - OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES)); + } + const pendingBytes = Math.min( + pending.length, + OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES - audio.length, + ); + return Buffer.concat( + [pending.subarray(pending.length - pendingBytes), audio], + pendingBytes + audio.length, + ); +} diff --git a/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts b/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts index 0c5e54d4d670..bcdf9afd2a7a 100644 --- a/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts +++ b/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts @@ -1,5 +1,9 @@ import { spawnSync } from "node:child_process"; import { describe, expect, it, vi } from "vitest"; +import { + OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES, + OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, +} from "./realtime-quicksilver-audio-buffer.js"; import { OpenAIQuicksilverGatewayBridge } from "./realtime-quicksilver-gateway-bridge.js"; import { OpenAIQuicksilverAudioPeer, @@ -46,7 +50,7 @@ type TestableAudioPeer = { }; peer: { connectionStateChange: { - execute(state: "closed" | "disconnected"): void; + execute(state: "closed" | "connected" | "disconnected"): void; }; }; transceiver: { @@ -574,66 +578,117 @@ describe("GPT-Live gateway relay bridge", () => { resolvePeer(); await connection; - expect(peer.sendAudio.mock.calls.map(([audio]) => audio)).toEqual([ - Buffer.from([0x7f, 0x41]), - Buffer.from([0x22, 0x23]), - ]); + expect(peer.sendAudio).toHaveBeenCalledWith(Buffer.from([0x7f, 0x41, 0x22, 0x23])); bridge.sendAudio(Buffer.from([0x30, 0x31])); - expect(peer.sendAudio).toHaveBeenCalledTimes(3); - } finally { - bridge.close(); - } - }); - - it("bounds queued microphone bytes before copying rejected audio", async () => { - const { bridge, connection, peer, resolvePeer } = createPendingPeerBridge(); - try { - bridge.sendAudio(Buffer.alloc(512 * 1024, 0x01)); - bridge.sendAudio(Buffer.alloc(512 * 1024, 0x02)); - const overflow = Buffer.alloc(1, 0x03); - const copy = vi.spyOn(Buffer, "from"); - try { - bridge.sendAudio(overflow); - expect(copy).not.toHaveBeenCalled(); - } finally { - copy.mockRestore(); - } - - resolvePeer(); - await connection; - expect(peer.sendAudio).toHaveBeenCalledTimes(2); - expect(peer.sendAudio.mock.calls.map(([audio]) => audio.byteLength)).toEqual([ - 512 * 1024, - 512 * 1024, - ]); } finally { bridge.close(); } }); - it("bounds queued microphone frame count before copying rejected audio", async () => { - const { bridge, connection, peer, resolvePeer } = createPendingPeerBridge(); + it("retains the newest five seconds through delayed peer adoption and the RTP pump", async () => { + vi.useFakeTimers(); + let peerCallbacks: + | Parameters[0]["callbacks"] + | undefined; + let resolvePeer: ((peer: OpenAIQuicksilverAudioPeerContract) => void) | undefined; + const peerPromise = new Promise((resolve) => { + resolvePeer = resolve; + }); + const onError = vi.fn(); + const bridge = new OpenAIQuicksilverGatewayBridge({ + providerConfig: {}, + model: "gpt-live-1-codex", + voice: "marin", + audioFormat: { encoding: "pcm16", sampleRateHz: 24_000, channels: 1 }, + onAudio: vi.fn(), + onClearAudio: vi.fn(), + onError, + runAgentConsult: vi.fn(async () => ({ text: "done" })), + logger: { debug: vi.fn(), warn: vi.fn() }, + resolveAuth: vi.fn(async () => ({ + type: "oauth" as const, + token: "oauth-token", + accountId: "account-1", + })), + createPeer: vi.fn((callbacks) => { + peerCallbacks = callbacks; + return peerPromise; + }), + fetchImpl: vi.fn(async () => createCallResponse("v=answer\r\n", "rtc_pending_audio")), + webSocketFactory: () => new FakeSocket(), + }); + const connection = bridge.connect(); + const source = Buffer.alloc(256 * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES); + for (let frame = 0; frame < 256; frame += 1) { + source + .subarray( + frame * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, + (frame + 1) * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, + ) + .fill(frame + 1); + } + for (let offset = 0; offset < source.length; offset += 8_192) { + const callerBuffer = Buffer.from(source.subarray(offset, offset + 8_192)); + bridge.sendAudio(callerBuffer); + callerBuffer.fill(0xff); + await vi.advanceTimersByTimeAsync(171); + } + if (!peerCallbacks) { + throw new Error("expected the bridge to start peer creation"); + } + const peer = await OpenAIQuicksilverAudioPeer.create({ + callbacks: peerCallbacks, + iceServers: [], + }); + vi.spyOn(peer, "createOffer").mockResolvedValue("v=offer\r\n"); + vi.spyOn(peer, "applyAnswer").mockResolvedValue(); + const testPeer = peer as unknown as TestableAudioPeer; + const sendRtp = vi + .spyOn(testPeer.state.transceiver.sender, "sendRtp") + .mockResolvedValue(undefined); + const emittedFrames: Buffer[] = []; + const takeNextRelayFrame = testPeer.takeNextRelayFrame.bind(testPeer); + vi.spyOn(testPeer, "takeNextRelayFrame").mockImplementation(() => { + const frame = takeNextRelayFrame(); + emittedFrames.push(frame); + return frame; + }); + const initialSequenceNumber = testPeer.sequenceNumber; + const initialTimestamp = testPeer.timestamp; try { - for (let index = 0; index < 320; index += 1) { - bridge.sendAudio(Buffer.alloc(2, index)); - } - const overflow = Buffer.alloc(2, 0xff); - const copy = vi.spyOn(Buffer, "from"); - try { - bridge.sendAudio(overflow); - expect(copy).not.toHaveBeenCalled(); - } finally { - copy.mockRestore(); - } - - resolvePeer(); + resolvePeer?.(peer); await connection; - expect(peer.sendAudio).toHaveBeenCalledTimes(320); - expect(peer.sendAudio.mock.calls.at(-1)?.[0]).toEqual(Buffer.alloc(2, 319)); + expect(testPeer.pendingAudio).toEqual( + source.subarray(source.length - OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES), + ); + + testPeer.state.peer.connectionStateChange.execute("connected"); + await vi.advanceTimersByTimeAsync(4_980); + + expect(testPeer.pendingAudio).toHaveLength(0); + expect(emittedFrames).toHaveLength(250); + expect(emittedFrames.map((frame) => frame[0])).toEqual([ + ...Array.from({ length: 249 }, (_, index) => index + 7), + 0, + ]); + expect(sendRtp).toHaveBeenCalledTimes(250); + const packets = sendRtp.mock.calls.map( + ([packet]) => + packet as { + header: { payloadType: number; sequenceNumber: number; timestamp: number }; + payload: Buffer; + }, + ); + expect(packets.every((packet) => packet.header.payloadType === 111)).toBe(true); + expect(packets.every((packet) => packet.payload.length > 0)).toBe(true); + expect(packets.at(-1)?.header.sequenceNumber).toBe((initialSequenceNumber + 249) & 0xffff); + expect(packets.at(-1)?.header.timestamp).toBe((initialTimestamp + 249 * 960) >>> 0); + expect(onError).not.toHaveBeenCalled(); } finally { bridge.close(); + vi.useRealTimers(); } }); @@ -653,17 +708,15 @@ describe("GPT-Live gateway relay bridge", () => { it("discards queued microphone audio when media peer creation fails", async () => { const { bridge, connection, peer, rejectPeer } = createPendingPeerBridge(); const pendingAudioState = bridge as unknown as { - pendingAudio: Buffer[]; - pendingAudioBytes: number; + pendingAudio: Buffer; }; bridge.sendAudio(Buffer.from([0x41, 0x42])); rejectPeer(new Error("media peer unavailable")); await expect(connection).rejects.toThrow("media peer unavailable"); - expect(pendingAudioState.pendingAudio).toEqual([]); - expect(pendingAudioState.pendingAudioBytes).toBe(0); + expect(pendingAudioState.pendingAudio).toHaveLength(0); bridge.sendAudio(Buffer.from([0x43, 0x44])); - expect(pendingAudioState.pendingAudio).toEqual([]); + expect(pendingAudioState.pendingAudio).toHaveLength(0); expect(peer.sendAudio).not.toHaveBeenCalled(); }); diff --git a/extensions/openai/realtime-quicksilver-gateway-bridge.ts b/extensions/openai/realtime-quicksilver-gateway-bridge.ts index 9e34ed4927b2..83422c071124 100644 --- a/extensions/openai/realtime-quicksilver-gateway-bridge.ts +++ b/extensions/openai/realtime-quicksilver-gateway-bridge.ts @@ -7,6 +7,7 @@ import type { RealtimeVoiceBridgeCreateRequest, } from "openclaw/plugin-sdk/realtime-voice"; import WebSocket, { type RawData } from "ws"; +import { appendOpenAIQuicksilverPendingAudio } from "./realtime-quicksilver-audio-buffer.js"; import { buildOpenAIQuicksilverDelegationPrompt, type OpenAIQuicksilverTranscriptEntry, @@ -40,7 +41,6 @@ const RELAY_SAMPLE_RATE = 24_000; const QUICKSILVER_SESSION_TTL_MS = 30 * 60_000; const QUICKSILVER_CONNECT_TIMEOUT_MS = 30_000; const WEBSOCKET_OPEN = 1; -const MAX_PENDING_AUDIO = { chunks: 320, bytes: 1024 * 1024 }; function toError(error: unknown): Error { return error instanceof Error ? error : new Error(String(error)); @@ -165,8 +165,7 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { private closed = false; private closeNotified = false; private peer: OpenAIQuicksilverAudioPeerContract | undefined; - private pendingAudio: Buffer[] = []; - private pendingAudioBytes = 0; + private pendingAudio: Buffer = Buffer.alloc(0); private ready = false; private sideband: ActiveSideband | undefined; private timer: ReturnType | undefined; @@ -186,15 +185,9 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { sendAudio(audio: Buffer): void { if (this.peer) { this.peer.sendAudio(audio); - } else if ( - !this.closed && - !this.abortController.signal.aborted && - this.pendingAudio.length < MAX_PENDING_AUDIO.chunks && - this.pendingAudioBytes + audio.byteLength <= MAX_PENDING_AUDIO.bytes - ) { + } else if (!this.closed && !this.abortController.signal.aborted) { // Relay capture starts before asynchronous peer creation and may recycle its input buffers. - this.pendingAudio.push(Buffer.from(audio)); - this.pendingAudioBytes += audio.byteLength; + this.pendingAudio = appendOpenAIQuicksilverPendingAudio(this.pendingAudio, audio); } } @@ -269,10 +262,10 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { () => undefined, ); this.peer = await waitForConnectStep(peerPromise, connectSignal); - for (const audio of this.pendingAudio.splice(0)) { - this.peer.sendAudio(audio); + if (this.pendingAudio.length > 0) { + this.peer.sendAudio(this.pendingAudio); + this.pendingAudio = Buffer.alloc(0); } - this.pendingAudioBytes = 0; const offerSdp = await waitForConnectStep(this.peer.createOffer(), connectSignal); const auth = await waitForConnectStep(this.config.resolveAuth(), connectSignal); const requestIds = { @@ -553,8 +546,7 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { private releaseResources(): void { releaseOpenAIQuicksilverSession(this); this.connected = false; - this.pendingAudio = []; - this.pendingAudioBytes = 0; + this.pendingAudio = Buffer.alloc(0); this.abortController.abort(new Error("GPT-Live gateway relay bridge closed")); this.consultController?.abort(new Error("GPT-Live delegation stopped")); this.consultController = undefined; diff --git a/extensions/openai/realtime-quicksilver-peer.runtime.ts b/extensions/openai/realtime-quicksilver-peer.runtime.ts index cd14cb012e3d..9eee52524c1a 100644 --- a/extensions/openai/realtime-quicksilver-peer.runtime.ts +++ b/extensions/openai/realtime-quicksilver-peer.runtime.ts @@ -1,15 +1,16 @@ // Lazy GPT-Live media runtime: werift peer plus WASM Opus framing and PCM conversion. import { randomInt } from "node:crypto"; import { resamplePcm } from "openclaw/plugin-sdk/realtime-voice"; +import { + appendOpenAIQuicksilverPendingAudio, + OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, +} from "./realtime-quicksilver-audio-buffer.js"; const QUICKSILVER_SAMPLE_RATE = 48_000; const RELAY_SAMPLE_RATE = 24_000; const QUICKSILVER_CHANNELS = 2; const OPUS_FRAME_SAMPLES = 960; const OPUS_FRAME_DURATION_MS = 20; -const RELAY_FRAME_SAMPLES = 480; -const RELAY_FRAME_BYTES = RELAY_FRAME_SAMPLES * 2; -const MAX_PENDING_RELAY_FRAMES = 250; const INBOUND_REORDER_DEPTH = 4; // More than two seconds behind cannot be useful 20 ms reordering; fail instead of corrupting Opus state. const INBOUND_MAX_LATE_PACKETS = 100; @@ -167,7 +168,7 @@ export class OpenAIQuicksilverAudioPeer implements OpenAIQuicksilverAudioPeerCon private activeInboundSsrc: number | undefined; private inboundRtpState: InboundRtpState = { pendingPackets: new Map() }; private mediaTimer: ReturnType | undefined; - private pendingAudio = Buffer.alloc(0); + private pendingAudio: Buffer = Buffer.alloc(0); private sequenceNumber = randomInt(0x1_0000); private subscribedTracks = new Set(); private timestamp = randomInt(0x1_0000_0000); @@ -222,16 +223,7 @@ export class OpenAIQuicksilverAudioPeer implements OpenAIQuicksilverAudioPeerCon if (this.closed || audio.length < 2) { return; } - const evenAudio = audio.subarray(0, audio.length - (audio.length % 2)); - this.pendingAudio = - this.pendingAudio.length > 0 - ? Buffer.concat([this.pendingAudio, evenAudio]) - : Buffer.from(evenAudio); - const maxPendingBytes = RELAY_FRAME_BYTES * MAX_PENDING_RELAY_FRAMES; - if (this.pendingAudio.length > maxPendingBytes) { - // Keep the newest complete frames. Old microphone audio is less useful than bounded latency. - this.pendingAudio = this.pendingAudio.subarray(this.pendingAudio.length - maxPendingBytes); - } + this.pendingAudio = appendOpenAIQuicksilverPendingAudio(this.pendingAudio, audio); } close(): void { @@ -440,8 +432,8 @@ export class OpenAIQuicksilverAudioPeer implements OpenAIQuicksilverAudioPeerCon private takeNextRelayFrame(): Buffer { // Relay ticks are framing boundaries: pad partial PCM now, or its tail survives // silence and is prepended to a later utterance as stale audio. - const frame = Buffer.alloc(RELAY_FRAME_BYTES); - const queuedBytes = Math.min(this.pendingAudio.length, RELAY_FRAME_BYTES); + const frame = Buffer.alloc(OPENAI_QUICKSILVER_RELAY_FRAME_BYTES); + const queuedBytes = Math.min(this.pendingAudio.length, OPENAI_QUICKSILVER_RELAY_FRAME_BYTES); if (queuedBytes > 0) { this.pendingAudio.copy(frame, 0, 0, queuedBytes); this.pendingAudio = this.pendingAudio.subarray(queuedBytes); From fd995d6f78ca389bd3238c5f5aa9e1bd7654fe81 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sat, 1 Aug 2026 13:32:17 +0800 Subject: [PATCH 3/5] test(openai): isolate gateway audio pipeline proof --- .../realtime-quicksilver-audio-buffer.test.ts | 12 +- .../realtime-quicksilver-audio-buffer.ts | 2 +- ...ealtime-quicksilver-audio-pipeline.test.ts | 142 ++++++++++++++++++ ...ealtime-quicksilver-gateway-bridge.test.ts | 110 -------------- 4 files changed, 150 insertions(+), 116 deletions(-) create mode 100644 extensions/openai/realtime-quicksilver-audio-pipeline.test.ts diff --git a/extensions/openai/realtime-quicksilver-audio-buffer.test.ts b/extensions/openai/realtime-quicksilver-audio-buffer.test.ts index af92af2dc172..85dfa89637e8 100644 --- a/extensions/openai/realtime-quicksilver-audio-buffer.test.ts +++ b/extensions/openai/realtime-quicksilver-audio-buffer.test.ts @@ -1,9 +1,11 @@ import { describe, expect, it } from "vitest"; import { appendOpenAIQuicksilverPendingAudio, - OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES, + OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, } from "./realtime-quicksilver-audio-buffer.js"; +const MAX_PENDING_AUDIO_BYTES = OPENAI_QUICKSILVER_RELAY_FRAME_BYTES * 250; + describe("GPT-Live pending microphone audio", () => { it("copies caller-owned PCM16 and drops an incomplete sample", () => { const source = Buffer.from([0x01, 0x02, 0x03]); @@ -21,16 +23,16 @@ describe("GPT-Live pending microphone audio", () => { }); it("retains the newest bounded tail across existing and oversized input", () => { - const existing = Buffer.alloc(OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES, 0x01); + const existing = Buffer.alloc(MAX_PENDING_AUDIO_BYTES, 0x01); const appended = appendOpenAIQuicksilverPendingAudio(existing, Buffer.from([0x02, 0x02])); - const oversized = Buffer.alloc(OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES + 2, 0x03); + const oversized = Buffer.alloc(MAX_PENDING_AUDIO_BYTES + 2, 0x03); - expect(appended).toHaveLength(OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES); + expect(appended).toHaveLength(MAX_PENDING_AUDIO_BYTES); expect(appended.subarray(0, -2).every((byte) => byte === 0x01)).toBe(true); expect(appended.subarray(-2)).toEqual(Buffer.from([0x02, 0x02])); expect( appendOpenAIQuicksilverPendingAudio(appended, oversized).equals( - Buffer.alloc(OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES, 0x03), + Buffer.alloc(MAX_PENDING_AUDIO_BYTES, 0x03), ), ).toBe(true); }); diff --git a/extensions/openai/realtime-quicksilver-audio-buffer.ts b/extensions/openai/realtime-quicksilver-audio-buffer.ts index 5616cb7635aa..144cb1916989 100644 --- a/extensions/openai/realtime-quicksilver-audio-buffer.ts +++ b/extensions/openai/realtime-quicksilver-audio-buffer.ts @@ -4,7 +4,7 @@ const MAX_PENDING_RELAY_FRAMES = 250; export const OPENAI_QUICKSILVER_RELAY_FRAME_BYTES = RELAY_FRAME_SAMPLES * 2; // One five-second tail spans peer startup and the connected media pump. // Keeping the newest PCM bounds latency without changing policy at adoption. -export const OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES = +const OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES = OPENAI_QUICKSILVER_RELAY_FRAME_BYTES * MAX_PENDING_RELAY_FRAMES; export function appendOpenAIQuicksilverPendingAudio(pending: Buffer, incoming: Buffer): Buffer { diff --git a/extensions/openai/realtime-quicksilver-audio-pipeline.test.ts b/extensions/openai/realtime-quicksilver-audio-pipeline.test.ts new file mode 100644 index 000000000000..69500a8cd915 --- /dev/null +++ b/extensions/openai/realtime-quicksilver-audio-pipeline.test.ts @@ -0,0 +1,142 @@ +import { describe, expect, it, vi } from "vitest"; +import { OPENAI_QUICKSILVER_RELAY_FRAME_BYTES } from "./realtime-quicksilver-audio-buffer.js"; +import { OpenAIQuicksilverGatewayBridge } from "./realtime-quicksilver-gateway-bridge.js"; +import { + OpenAIQuicksilverAudioPeer, + type OpenAIQuicksilverAudioPeerContract, +} from "./realtime-quicksilver-peer.runtime.js"; +import { createCallResponse, FakeSocket } from "./realtime-quicksilver.test-helpers.js"; + +const MAX_PENDING_RELAY_FRAMES = 250; +const MAX_PENDING_AUDIO_BYTES = OPENAI_QUICKSILVER_RELAY_FRAME_BYTES * MAX_PENDING_RELAY_FRAMES; + +type TestableAudioPeer = { + pendingAudio: Buffer; + sequenceNumber: number; + timestamp: number; + takeNextRelayFrame(): Buffer; + state: { + peer: { + connectionStateChange: { + execute(state: "connected"): void; + }; + }; + transceiver: { + sender: { + sendRtp(packet: unknown): Promise; + }; + }; + }; +}; + +describe("GPT-Live gateway microphone audio pipeline", () => { + it("retains the newest five seconds through delayed peer adoption and the RTP pump", async () => { + vi.useFakeTimers(); + let peerCallbacks: + | Parameters[0]["callbacks"] + | undefined; + let resolvePeer: ((peer: OpenAIQuicksilverAudioPeerContract) => void) | undefined; + const peerPromise = new Promise((resolve) => { + resolvePeer = resolve; + }); + const onError = vi.fn(); + const bridge = new OpenAIQuicksilverGatewayBridge({ + providerConfig: {}, + model: "gpt-live-1-codex", + voice: "marin", + audioFormat: { encoding: "pcm16", sampleRateHz: 24_000, channels: 1 }, + onAudio: vi.fn(), + onClearAudio: vi.fn(), + onError, + runAgentConsult: vi.fn(async () => ({ text: "done" })), + logger: { debug: vi.fn(), warn: vi.fn() }, + resolveAuth: vi.fn(async () => ({ + type: "oauth" as const, + token: "oauth-token", + accountId: "account-1", + })), + createPeer: vi.fn((callbacks) => { + peerCallbacks = callbacks; + return peerPromise; + }), + fetchImpl: vi.fn(async () => createCallResponse("v=answer\r\n", "rtc_pending_audio")), + webSocketFactory: () => new FakeSocket(), + }); + const connection = bridge.connect(); + const source = Buffer.alloc(256 * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES); + for (let frame = 0; frame < 256; frame += 1) { + source + .subarray( + frame * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, + (frame + 1) * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, + ) + .fill(frame + 1); + } + for (let offset = 0; offset < source.length; offset += 8_192) { + const callerBuffer = Buffer.from(source.subarray(offset, offset + 8_192)); + bridge.sendAudio(callerBuffer); + callerBuffer.fill(0xff); + await vi.advanceTimersByTimeAsync(171); + } + if (!peerCallbacks) { + throw new Error("expected the bridge to start peer creation"); + } + const peer = await OpenAIQuicksilverAudioPeer.create({ + callbacks: peerCallbacks, + iceServers: [], + }); + vi.spyOn(peer, "createOffer").mockResolvedValue("v=offer\r\n"); + vi.spyOn(peer, "applyAnswer").mockResolvedValue(); + const testPeer = peer as unknown as TestableAudioPeer; + const sendRtp = vi + .spyOn(testPeer.state.transceiver.sender, "sendRtp") + .mockResolvedValue(undefined); + const emittedFrames: Buffer[] = []; + const takeNextRelayFrame = testPeer.takeNextRelayFrame.bind(testPeer); + vi.spyOn(testPeer, "takeNextRelayFrame").mockImplementation(() => { + const frame = takeNextRelayFrame(); + emittedFrames.push(frame); + return frame; + }); + const initialSequenceNumber = testPeer.sequenceNumber; + const initialTimestamp = testPeer.timestamp; + try { + resolvePeer?.(peer); + await connection; + + expect(testPeer.pendingAudio).toEqual( + source.subarray(source.length - MAX_PENDING_AUDIO_BYTES), + ); + + testPeer.state.peer.connectionStateChange.execute("connected"); + await vi.advanceTimersByTimeAsync(4_980); + + expect(testPeer.pendingAudio).toHaveLength(0); + expect(emittedFrames).toHaveLength(MAX_PENDING_RELAY_FRAMES); + expect(emittedFrames.map((frame) => frame[0])).toEqual([ + ...Array.from({ length: MAX_PENDING_RELAY_FRAMES - 1 }, (_, index) => index + 7), + 0, + ]); + expect(sendRtp).toHaveBeenCalledTimes(MAX_PENDING_RELAY_FRAMES); + const packets = sendRtp.mock.calls.map( + ([packet]) => + packet as { + header: { payloadType: number; sequenceNumber: number; timestamp: number }; + payload: Buffer; + }, + ); + expect(packets.every((packet) => packet.header.payloadType === 111)).toBe(true); + expect(packets.every((packet) => packet.payload.length > 0)).toBe(true); + expect(packets.at(-1)?.header.sequenceNumber).toBe( + (initialSequenceNumber + MAX_PENDING_RELAY_FRAMES - 1) & 0xffff, + ); + expect(packets.at(-1)?.header.timestamp).toBe( + (initialTimestamp + (MAX_PENDING_RELAY_FRAMES - 1) * 960) >>> 0, + ); + expect(onError).not.toHaveBeenCalled(); + } finally { + bridge.close(); + vi.useRealTimers(); + } + }); +}); diff --git a/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts b/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts index bcdf9afd2a7a..3bc00a66c2d2 100644 --- a/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts +++ b/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts @@ -1,9 +1,5 @@ import { spawnSync } from "node:child_process"; import { describe, expect, it, vi } from "vitest"; -import { - OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES, - OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, -} from "./realtime-quicksilver-audio-buffer.js"; import { OpenAIQuicksilverGatewayBridge } from "./realtime-quicksilver-gateway-bridge.js"; import { OpenAIQuicksilverAudioPeer, @@ -586,112 +582,6 @@ describe("GPT-Live gateway relay bridge", () => { } }); - it("retains the newest five seconds through delayed peer adoption and the RTP pump", async () => { - vi.useFakeTimers(); - let peerCallbacks: - | Parameters[0]["callbacks"] - | undefined; - let resolvePeer: ((peer: OpenAIQuicksilverAudioPeerContract) => void) | undefined; - const peerPromise = new Promise((resolve) => { - resolvePeer = resolve; - }); - const onError = vi.fn(); - const bridge = new OpenAIQuicksilverGatewayBridge({ - providerConfig: {}, - model: "gpt-live-1-codex", - voice: "marin", - audioFormat: { encoding: "pcm16", sampleRateHz: 24_000, channels: 1 }, - onAudio: vi.fn(), - onClearAudio: vi.fn(), - onError, - runAgentConsult: vi.fn(async () => ({ text: "done" })), - logger: { debug: vi.fn(), warn: vi.fn() }, - resolveAuth: vi.fn(async () => ({ - type: "oauth" as const, - token: "oauth-token", - accountId: "account-1", - })), - createPeer: vi.fn((callbacks) => { - peerCallbacks = callbacks; - return peerPromise; - }), - fetchImpl: vi.fn(async () => createCallResponse("v=answer\r\n", "rtc_pending_audio")), - webSocketFactory: () => new FakeSocket(), - }); - const connection = bridge.connect(); - const source = Buffer.alloc(256 * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES); - for (let frame = 0; frame < 256; frame += 1) { - source - .subarray( - frame * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, - (frame + 1) * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, - ) - .fill(frame + 1); - } - for (let offset = 0; offset < source.length; offset += 8_192) { - const callerBuffer = Buffer.from(source.subarray(offset, offset + 8_192)); - bridge.sendAudio(callerBuffer); - callerBuffer.fill(0xff); - await vi.advanceTimersByTimeAsync(171); - } - if (!peerCallbacks) { - throw new Error("expected the bridge to start peer creation"); - } - const peer = await OpenAIQuicksilverAudioPeer.create({ - callbacks: peerCallbacks, - iceServers: [], - }); - vi.spyOn(peer, "createOffer").mockResolvedValue("v=offer\r\n"); - vi.spyOn(peer, "applyAnswer").mockResolvedValue(); - const testPeer = peer as unknown as TestableAudioPeer; - const sendRtp = vi - .spyOn(testPeer.state.transceiver.sender, "sendRtp") - .mockResolvedValue(undefined); - const emittedFrames: Buffer[] = []; - const takeNextRelayFrame = testPeer.takeNextRelayFrame.bind(testPeer); - vi.spyOn(testPeer, "takeNextRelayFrame").mockImplementation(() => { - const frame = takeNextRelayFrame(); - emittedFrames.push(frame); - return frame; - }); - const initialSequenceNumber = testPeer.sequenceNumber; - const initialTimestamp = testPeer.timestamp; - try { - resolvePeer?.(peer); - await connection; - - expect(testPeer.pendingAudio).toEqual( - source.subarray(source.length - OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES), - ); - - testPeer.state.peer.connectionStateChange.execute("connected"); - await vi.advanceTimersByTimeAsync(4_980); - - expect(testPeer.pendingAudio).toHaveLength(0); - expect(emittedFrames).toHaveLength(250); - expect(emittedFrames.map((frame) => frame[0])).toEqual([ - ...Array.from({ length: 249 }, (_, index) => index + 7), - 0, - ]); - expect(sendRtp).toHaveBeenCalledTimes(250); - const packets = sendRtp.mock.calls.map( - ([packet]) => - packet as { - header: { payloadType: number; sequenceNumber: number; timestamp: number }; - payload: Buffer; - }, - ); - expect(packets.every((packet) => packet.header.payloadType === 111)).toBe(true); - expect(packets.every((packet) => packet.payload.length > 0)).toBe(true); - expect(packets.at(-1)?.header.sequenceNumber).toBe((initialSequenceNumber + 249) & 0xffff); - expect(packets.at(-1)?.header.timestamp).toBe((initialTimestamp + 249 * 960) >>> 0); - expect(onError).not.toHaveBeenCalled(); - } finally { - bridge.close(); - vi.useRealTimers(); - } - }); - it("discards queued microphone audio when closed before the media peer resolves", async () => { const { bridge, connection, peer, resolvePeer } = createPendingPeerBridge(); bridge.sendAudio(Buffer.from([0x41, 0x42])); From da94915d064989bf67455c1bd5938e58c42ba8d0 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sun, 2 Aug 2026 03:07:01 +0800 Subject: [PATCH 4/5] test(openai): prove gateway startup audio delivery --- .../realtime-quicksilver-audio-buffer.test.ts | 14 +- ...ealtime-quicksilver-audio-pipeline.test.ts | 27 +-- ...me-quicksilver-gateway-bridge.live.test.ts | 161 ++++++++++++++++++ ...ealtime-quicksilver-gateway-bridge.test.ts | 36 +++- 4 files changed, 221 insertions(+), 17 deletions(-) diff --git a/extensions/openai/realtime-quicksilver-audio-buffer.test.ts b/extensions/openai/realtime-quicksilver-audio-buffer.test.ts index 85dfa89637e8..a43667f38198 100644 --- a/extensions/openai/realtime-quicksilver-audio-buffer.test.ts +++ b/extensions/openai/realtime-quicksilver-audio-buffer.test.ts @@ -25,15 +25,17 @@ describe("GPT-Live pending microphone audio", () => { it("retains the newest bounded tail across existing and oversized input", () => { const existing = Buffer.alloc(MAX_PENDING_AUDIO_BYTES, 0x01); const appended = appendOpenAIQuicksilverPendingAudio(existing, Buffer.from([0x02, 0x02])); - const oversized = Buffer.alloc(MAX_PENDING_AUDIO_BYTES + 2, 0x03); + const oversized = Buffer.alloc(MAX_PENDING_AUDIO_BYTES + 4, 0x03); + oversized.writeUInt16LE(0x1111, 0); + oversized.writeUInt16LE(0x2222, oversized.length - 2); + const expectedOversizedTail = Buffer.from(oversized.subarray(4)); + const oversizedResult = appendOpenAIQuicksilverPendingAudio(appended, oversized); + oversized.fill(0xff); expect(appended).toHaveLength(MAX_PENDING_AUDIO_BYTES); expect(appended.subarray(0, -2).every((byte) => byte === 0x01)).toBe(true); expect(appended.subarray(-2)).toEqual(Buffer.from([0x02, 0x02])); - expect( - appendOpenAIQuicksilverPendingAudio(appended, oversized).equals( - Buffer.alloc(MAX_PENDING_AUDIO_BYTES, 0x03), - ), - ).toBe(true); + expect(oversizedResult).toEqual(expectedOversizedTail); + expect(oversizedResult.readUInt16LE(-2 + oversizedResult.length)).toBe(0x2222); }); }); diff --git a/extensions/openai/realtime-quicksilver-audio-pipeline.test.ts b/extensions/openai/realtime-quicksilver-audio-pipeline.test.ts index 69500a8cd915..2061070280ad 100644 --- a/extensions/openai/realtime-quicksilver-audio-pipeline.test.ts +++ b/extensions/openai/realtime-quicksilver-audio-pipeline.test.ts @@ -29,6 +29,10 @@ type TestableAudioPeer = { }; }; +type TestableGatewayBridge = { + pendingAudio: Buffer; +}; + describe("GPT-Live gateway microphone audio pipeline", () => { it("retains the newest five seconds through delayed peer adoption and the RTP pump", async () => { vi.useFakeTimers(); @@ -65,12 +69,13 @@ describe("GPT-Live gateway microphone audio pipeline", () => { const connection = bridge.connect(); const source = Buffer.alloc(256 * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES); for (let frame = 0; frame < 256; frame += 1) { - source - .subarray( - frame * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, - (frame + 1) * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, - ) - .fill(frame + 1); + const frameBuffer = source.subarray( + frame * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, + (frame + 1) * OPENAI_QUICKSILVER_RELAY_FRAME_BYTES, + ); + for (let sample = 0; sample < frameBuffer.length; sample += 2) { + frameBuffer.writeUInt16LE(frame + 1, sample); + } } for (let offset = 0; offset < source.length; offset += 8_192) { const callerBuffer = Buffer.from(source.subarray(offset, offset + 8_192)); @@ -78,6 +83,9 @@ describe("GPT-Live gateway microphone audio pipeline", () => { callerBuffer.fill(0xff); await vi.advanceTimersByTimeAsync(171); } + expect((bridge as unknown as TestableGatewayBridge).pendingAudio).toEqual( + source.subarray(source.length - MAX_PENDING_AUDIO_BYTES), + ); if (!peerCallbacks) { throw new Error("expected the bridge to start peer creation"); } @@ -113,10 +121,9 @@ describe("GPT-Live gateway microphone audio pipeline", () => { expect(testPeer.pendingAudio).toHaveLength(0); expect(emittedFrames).toHaveLength(MAX_PENDING_RELAY_FRAMES); - expect(emittedFrames.map((frame) => frame[0])).toEqual([ - ...Array.from({ length: MAX_PENDING_RELAY_FRAMES - 1 }, (_, index) => index + 7), - 0, - ]); + expect(emittedFrames.map((frame) => frame.readUInt16LE(0))).toEqual( + Array.from({ length: MAX_PENDING_RELAY_FRAMES }, (_, index) => index + 7), + ); expect(sendRtp).toHaveBeenCalledTimes(MAX_PENDING_RELAY_FRAMES); const packets = sendRtp.mock.calls.map( ([packet]) => diff --git a/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts b/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts index 41b542444519..0174dde16b67 100644 --- a/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts +++ b/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts @@ -2,13 +2,38 @@ import { readCodexCliCredentialsCached } from "openclaw/plugin-sdk/provider-auth import { describe, expect, it } from "vitest"; import { resolveCodexAuthIdentity } from "./openai-chatgpt-auth-identity.js"; import { OpenAIQuicksilverGatewayBridge } from "./realtime-quicksilver-gateway-bridge.js"; +import { + OpenAIQuicksilverAudioPeer, + type OpenAIQuicksilverAudioPeerContract, +} from "./realtime-quicksilver-peer.runtime.js"; import { resolveOpenAIChatGptSubscriptionAuth } from "./realtime-quicksilver-session.js"; import type { OpenAIQuicksilverAuth } from "./realtime-quicksilver-wire.js"; +import { buildOpenAISpeechProvider } from "./speech-provider.js"; const LIVE_ENABLED = process.env.OPENCLAW_LIVE_TEST === "1" && process.env.OPENCLAW_LIVE_GPT_LIVE === "1"; const describeLive = LIVE_ENABLED ? describe : describe.skip; const LIVE_TIMEOUT_MS = 60_000; +const MAX_PENDING_AUDIO_BYTES = 240_000; + +type TestableGatewayBridge = { + pendingAudio: Buffer; +}; + +async function waitForLiveCondition( + predicate: () => boolean, + describeFailure: () => string, + timeoutMs = 45_000, +): Promise { + const startedAt = Date.now(); + while (Date.now() - startedAt < timeoutMs) { + if (predicate()) { + return; + } + await new Promise((resolve) => setTimeout(resolve, 100)); + } + throw new Error(describeFailure()); +} async function resolveLiveOAuthProfile(): Promise< Extract | undefined @@ -86,4 +111,140 @@ describeLive("OpenAI GPT-Live gateway WebRTC peer", () => { }, LIVE_TIMEOUT_MS, ); + + it( + "delivers microphone speech queued before the real media peer is adopted", + async ({ skip }) => { + const apiKey = process.env.OPENAI_API_KEY?.trim(); + if (!apiKey) { + skip("No OpenAI Platform API key is available for the speech fixture"); + return; + } + + const speechProvider = buildOpenAISpeechProvider(); + const synthesized = await speechProvider.synthesizeTelephony?.({ + text: "Please delegate the word glacier.", + cfg: { plugins: { enabled: true } } as never, + providerConfig: { + apiKey, + baseUrl: "https://api.openai.com/v1", + model: "gpt-4o-mini-tts", + voice: "alloy", + speed: 1.4, + }, + timeoutMs: 45_000, + }); + if (!synthesized) { + throw new Error("OpenAI speech provider did not return a telephony fixture"); + } + expect(synthesized.outputFormat).toBe("pcm"); + expect(synthesized.sampleRate).toBe(24_000); + const inputAudio = Buffer.concat([synthesized.audioBuffer, Buffer.alloc(24_000 * 2)]); + expect(inputAudio.byteLength).toBeLessThanOrEqual(MAX_PENDING_AUDIO_BYTES); + + let releasePeerAdoption!: () => void; + let peerCreated!: () => void; + const peerAdoption = new Promise((resolve) => { + releasePeerAdoption = resolve; + }); + const peerCreation = new Promise((resolve) => { + peerCreated = resolve; + }); + const eventTypes: string[] = []; + const finalUserTranscripts: string[] = []; + const errors: Error[] = []; + let closeNotifications = 0; + let closed = false; + let lateAudioBytes = 0; + const bridge = new OpenAIQuicksilverGatewayBridge({ + providerConfig: {}, + model: "gpt-live-1-boulder-alpha", + voice: "marin", + instructions: "Listen to the user. Do not speak or delegate.", + audioFormat: { encoding: "pcm16", sampleRateHz: 24_000, channels: 1 }, + onAudio: (audio) => { + if (closed) { + lateAudioBytes += audio.length; + } + }, + onClearAudio: () => undefined, + onEvent: (event) => eventTypes.push(event.type), + onReady: () => undefined, + onTranscript: (role, text, final) => { + if (role === "user" && final) { + finalUserTranscripts.push(text); + } + }, + onClose: () => { + closeNotifications += 1; + }, + onError: (error) => errors.push(error), + runAgentConsult: async () => ({ text: "Unexpected delegation." }), + logger: { debug: () => undefined, warn: () => undefined }, + resolveAuth: async () => ({ type: "api-key", token: apiKey }), + createPeer: async (callbacks, signal): Promise => { + const peer = await OpenAIQuicksilverAudioPeer.create({ callbacks, signal }); + peerCreated(); + await peerAdoption; + return peer; + }, + }); + const testBridge = bridge as unknown as TestableGatewayBridge; + + try { + const connection = bridge.connect(); + await Promise.race([ + peerCreation, + connection.then(() => { + throw new Error("Gateway bridge connected before the media peer adoption gate"); + }), + ]); + for (let offset = 0; offset < inputAudio.length; offset += 8_192) { + bridge.sendAudio(Buffer.from(inputAudio.subarray(offset, offset + 8_192))); + } + const prePeerPendingBytes = testBridge.pendingAudio.length; + expect(prePeerPendingBytes).toBe(inputAudio.length); + + releasePeerAdoption(); + await connection; + await waitForLiveCondition( + () => finalUserTranscripts.some((text) => text.toLowerCase().includes("glacier")), + () => + `GPT-Live did not transcribe startup audio: transcripts=${finalUserTranscripts.length} errors=${errors.map((error) => error.message).join(";")}`, + 30_000, + ); + + const postAdoptionPendingBytes = testBridge.pendingAudio.length; + expect(postAdoptionPendingBytes).toBe(0); + expect(eventTypes).toContain("turn.done"); + expect(bridge.isConnected()).toBe(true); + expect(errors).toStrictEqual([]); + + closed = true; + bridge.close(); + bridge.close(); + await new Promise((resolve) => setTimeout(resolve, 250)); + + expect(closeNotifications).toBe(1); + expect(lateAudioBytes).toBe(0); + expect(errors).toStrictEqual([]); + console.log( + JSON.stringify({ + proof: "gpt-live-gateway-pre-peer-transcription", + prePeerPendingBytes, + postAdoptionPendingBytes, + userTranscriptMarker: true, + closeNotifications, + lateAudioBytes, + errors: errors.length, + result: "pass", + }), + ); + } finally { + releasePeerAdoption(); + bridge.close(); + } + }, + LIVE_TIMEOUT_MS, + ); }); diff --git a/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts b/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts index 3bc00a66c2d2..839363f3a345 100644 --- a/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts +++ b/extensions/openai/realtime-quicksilver-gateway-bridge.test.ts @@ -1,5 +1,6 @@ import { spawnSync } from "node:child_process"; import { describe, expect, it, vi } from "vitest"; +import { OPENAI_QUICKSILVER_RELAY_FRAME_BYTES } from "./realtime-quicksilver-audio-buffer.js"; import { OpenAIQuicksilverGatewayBridge } from "./realtime-quicksilver-gateway-bridge.js"; import { OpenAIQuicksilverAudioPeer, @@ -58,6 +59,7 @@ type TestableAudioPeer = { }; type TestableGatewayBridge = { + pendingAudio: Buffer; sideband?: { socket: FakeSocket; requestIds: { realtimeSessionId: string; sessionId: string; threadId: string }; @@ -357,6 +359,28 @@ describe("GPT-Live werift audio peer", () => { } }); + it("retains only the newest five seconds and releases it on close", async () => { + const peer = await OpenAIQuicksilverAudioPeer.create({ + callbacks: { onAudio: vi.fn(), onError: vi.fn() }, + iceServers: [], + }); + const testPeer = peer as unknown as TestableAudioPeer; + const maxPendingAudioBytes = OPENAI_QUICKSILVER_RELAY_FRAME_BYTES * 250; + const source = Buffer.alloc(maxPendingAudioBytes + OPENAI_QUICKSILVER_RELAY_FRAME_BYTES); + source.fill(0x11, 0, OPENAI_QUICKSILVER_RELAY_FRAME_BYTES); + source.fill(0x22, OPENAI_QUICKSILVER_RELAY_FRAME_BYTES); + const expectedTail = Buffer.from(source.subarray(OPENAI_QUICKSILVER_RELAY_FRAME_BYTES)); + + peer.sendAudio(source); + source.fill(0xff); + expect(testPeer.pendingAudio).toEqual(expectedTail); + + peer.close(); + expect(testPeer.pendingAudio).toHaveLength(0); + peer.sendAudio(Buffer.from([0x01, 0x02])); + expect(testPeer.pendingAudio).toHaveLength(0); + }); + it("consumes and zero-pads a sub-frame audio tail on the next tick", async () => { const peer = await OpenAIQuicksilverAudioPeer.create({ callbacks: { onAudio: vi.fn(), onError: vi.fn() }, @@ -535,6 +559,7 @@ describe("GPT-Live gateway relay bridge", () => { sendAudio: vi.fn(), close: vi.fn(), } satisfies OpenAIQuicksilverAudioPeerContract; + const onClose = vi.fn(); const bridge = new OpenAIQuicksilverGatewayBridge({ providerConfig: {}, model: "gpt-live-1-codex", @@ -542,6 +567,7 @@ describe("GPT-Live gateway relay bridge", () => { audioFormat: { encoding: "pcm16", sampleRateHz: 24_000, channels: 1 }, onAudio: vi.fn(), onClearAudio: vi.fn(), + onClose, runAgentConsult: vi.fn(async () => ({ text: "done" })), logger: { debug: vi.fn(), warn: vi.fn() }, resolveAuth: vi.fn(async () => ({ @@ -557,6 +583,7 @@ describe("GPT-Live gateway relay bridge", () => { return { bridge, connection, + onClose, peer, rejectPeer: (error: Error) => rejectPeer?.(error), resolvePeer: () => resolvePeer?.(peer), @@ -566,6 +593,7 @@ describe("GPT-Live gateway relay bridge", () => { it("preserves caller-owned microphone frames while the media peer is starting", async () => { const { bridge, connection, peer, resolvePeer } = createPendingPeerBridge(); try { + expect(bridge.connect()).toBe(connection); const source = Buffer.from([0x7f, 0x41]); bridge.sendAudio(source); source.fill(0); @@ -583,9 +611,15 @@ describe("GPT-Live gateway relay bridge", () => { }); it("discards queued microphone audio when closed before the media peer resolves", async () => { - const { bridge, connection, peer, resolvePeer } = createPendingPeerBridge(); + const { bridge, connection, onClose, peer, resolvePeer } = createPendingPeerBridge(); + const testBridge = bridge as unknown as TestableGatewayBridge; bridge.sendAudio(Buffer.from([0x41, 0x42])); bridge.close(); + bridge.close(); + + expect(testBridge.pendingAudio).toHaveLength(0); + expect(onClose).toHaveBeenCalledOnce(); + expect(onClose).toHaveBeenCalledWith("completed"); resolvePeer(); await expect(connection).rejects.toThrow("GPT-Live gateway relay bridge closed"); From 75fc97369066aa940e9c1ad50ce7556eb5625baf Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Sun, 2 Aug 2026 03:20:30 +0800 Subject: [PATCH 5/5] test(openai): satisfy live proof lint --- .../realtime-quicksilver-gateway-bridge.live.test.ts | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) diff --git a/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts b/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts index 0174dde16b67..3e7cac975661 100644 --- a/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts +++ b/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts @@ -30,7 +30,9 @@ async function waitForLiveCondition( if (predicate()) { return; } - await new Promise((resolve) => setTimeout(resolve, 100)); + await new Promise((resolve) => { + setTimeout(resolve, 100); + }); } throw new Error(describeFailure()); } @@ -223,7 +225,9 @@ describeLive("OpenAI GPT-Live gateway WebRTC peer", () => { closed = true; bridge.close(); bridge.close(); - await new Promise((resolve) => setTimeout(resolve, 250)); + await new Promise((resolve) => { + setTimeout(resolve, 250); + }); expect(closeNotifications).toBe(1); expect(lateAudioBytes).toBe(0);