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..a43667f38198 --- /dev/null +++ b/extensions/openai/realtime-quicksilver-audio-buffer.test.ts @@ -0,0 +1,41 @@ +import { describe, expect, it } from "vitest"; +import { + appendOpenAIQuicksilverPendingAudio, + 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]); + 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(MAX_PENDING_AUDIO_BYTES, 0x01); + const appended = appendOpenAIQuicksilverPendingAudio(existing, Buffer.from([0x02, 0x02])); + 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(oversizedResult).toEqual(expectedOversizedTail); + expect(oversizedResult.readUInt16LE(-2 + oversizedResult.length)).toBe(0x2222); + }); +}); diff --git a/extensions/openai/realtime-quicksilver-audio-buffer.ts b/extensions/openai/realtime-quicksilver-audio-buffer.ts new file mode 100644 index 000000000000..144cb1916989 --- /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. +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-audio-pipeline.test.ts b/extensions/openai/realtime-quicksilver-audio-pipeline.test.ts new file mode 100644 index 000000000000..2061070280ad --- /dev/null +++ b/extensions/openai/realtime-quicksilver-audio-pipeline.test.ts @@ -0,0 +1,149 @@ +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; + }; + }; + }; +}; + +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(); + 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) { + 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)); + bridge.sendAudio(callerBuffer); + 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"); + } + 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.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]) => + 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.live.test.ts b/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts index 41b542444519..3e7cac975661 100644 --- a/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts +++ b/extensions/openai/realtime-quicksilver-gateway-bridge.live.test.ts @@ -2,13 +2,40 @@ 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 +113,142 @@ 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 c35b58dd6392..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, @@ -46,7 +47,7 @@ type TestableAudioPeer = { }; peer: { connectionStateChange: { - execute(state: "closed" | "disconnected"): void; + execute(state: "closed" | "connected" | "disconnected"): void; }; }; transceiver: { @@ -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() }, @@ -522,6 +546,104 @@ 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 onClose = 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(), + onClose, + 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, + onClose, + 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 { + expect(bridge.connect()).toBe(connection); + const source = Buffer.from([0x7f, 0x41]); + bridge.sendAudio(source); + source.fill(0); + bridge.sendAudio(Buffer.from([0x22, 0x23])); + + resolvePeer(); + await connection; + + expect(peer.sendAudio).toHaveBeenCalledWith(Buffer.from([0x7f, 0x41, 0x22, 0x23])); + bridge.sendAudio(Buffer.from([0x30, 0x31])); + expect(peer.sendAudio).toHaveBeenCalledTimes(2); + } finally { + bridge.close(); + } + }); + + it("discards queued microphone audio when closed before the media peer resolves", async () => { + 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"); + 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; + }; + bridge.sendAudio(Buffer.from([0x41, 0x42])); + rejectPeer(new Error("media peer unavailable")); + + await expect(connection).rejects.toThrow("media peer unavailable"); + expect(pendingAudioState.pendingAudio).toHaveLength(0); + bridge.sendAudio(Buffer.from([0x43, 0x44])); + expect(pendingAudioState.pendingAudio).toHaveLength(0); + 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..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, @@ -164,6 +165,7 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { private closed = false; private closeNotified = false; private peer: OpenAIQuicksilverAudioPeerContract | undefined; + private pendingAudio: Buffer = Buffer.alloc(0); private ready = false; private sideband: ActiveSideband | undefined; private timer: ReturnType | undefined; @@ -181,7 +183,12 @@ 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) { + // Relay capture starts before asynchronous peer creation and may recycle its input buffers. + this.pendingAudio = appendOpenAIQuicksilverPendingAudio(this.pendingAudio, audio); + } } setMediaTimestamp(_ts: number): void {} @@ -255,6 +262,10 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { () => undefined, ); this.peer = await waitForConnectStep(peerPromise, connectSignal); + if (this.pendingAudio.length > 0) { + this.peer.sendAudio(this.pendingAudio); + this.pendingAudio = Buffer.alloc(0); + } const offerSdp = await waitForConnectStep(this.peer.createOffer(), connectSignal); const auth = await waitForConnectStep(this.config.resolveAuth(), connectSignal); const requestIds = { @@ -535,6 +546,7 @@ export class OpenAIQuicksilverGatewayBridge implements RealtimeVoiceBridge { private releaseResources(): void { releaseOpenAIQuicksilverSession(this); this.connected = false; + 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);