mirror of
https://github.com/openclaw/openclaw.git
synced 2026-08-02 05:01:35 +00:00
Merge repaired PR #117057 onto current main
* commit '657450148cae1e3eb3f05e9fcb709ca60ff8d71d': test(openai): satisfy live proof lint test(openai): prove gateway startup audio delivery test(openai): isolate gateway audio pipeline proof fix(openai): unify gateway microphone buffering fix(openai): preserve gateway microphone startup audio
This commit is contained in:
41
extensions/openai/realtime-quicksilver-audio-buffer.test.ts
Normal file
41
extensions/openai/realtime-quicksilver-audio-buffer.test.ts
Normal file
@@ -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);
|
||||
});
|
||||
});
|
||||
27
extensions/openai/realtime-quicksilver-audio-buffer.ts
Normal file
27
extensions/openai/realtime-quicksilver-audio-buffer.ts
Normal file
@@ -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,
|
||||
);
|
||||
}
|
||||
149
extensions/openai/realtime-quicksilver-audio-pipeline.test.ts
Normal file
149
extensions/openai/realtime-quicksilver-audio-pipeline.test.ts
Normal file
@@ -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<void>;
|
||||
};
|
||||
};
|
||||
};
|
||||
};
|
||||
|
||||
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<typeof OpenAIQuicksilverAudioPeer.create>[0]["callbacks"]
|
||||
| undefined;
|
||||
let resolvePeer: ((peer: OpenAIQuicksilverAudioPeerContract) => void) | undefined;
|
||||
const peerPromise = new Promise<OpenAIQuicksilverAudioPeerContract>((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();
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -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<void> {
|
||||
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<OpenAIQuicksilverAuth, { type: "oauth" }> | 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<void>((resolve) => {
|
||||
releasePeerAdoption = resolve;
|
||||
});
|
||||
const peerCreation = new Promise<void>((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<OpenAIQuicksilverAudioPeerContract> => {
|
||||
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,
|
||||
);
|
||||
});
|
||||
|
||||
@@ -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<OpenAIQuicksilverAudioPeerContract>((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");
|
||||
|
||||
@@ -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<typeof setTimeout> | 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;
|
||||
|
||||
@@ -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<typeof setInterval> | undefined;
|
||||
private pendingAudio = Buffer.alloc(0);
|
||||
private pendingAudio: Buffer = Buffer.alloc(0);
|
||||
private sequenceNumber = randomInt(0x1_0000);
|
||||
private subscribedTracks = new Set<string>();
|
||||
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);
|
||||
|
||||
Reference in New Issue
Block a user