openclaw/extensions/openai/realtime-quicksilver-audio-buffer.ts
RoboClaw 41ee7fb242
fix(voice): keep audio flowing while the Gateway is busy (#154119)
* fix(voice): keep audio flowing while the Gateway is busy

Move Discord voice transport, packet pacing, codecs and capture deadlines into a session-owned worker. Run GPT Live WebRTC and WebSocket media in workers and connect continuous playback with a bounded direct audio port. Preserve authorization, selected-agent work, recording receipts and transcript ownership on main, with generation fencing and ordered teardown.

Validated affected tests, the capture-finalization red/green regression, source and compiled worker lifecycle, root and plugin builds, type/lint/format and line-cap checks.

Co-authored-by: steipete <58493+steipete@users.noreply.github.com>

* fix(voice): keep audio flowing while the Gateway is busy

Worked on by:
- @steipete

Co-authored-by: steipete <58493+steipete@users.noreply.github.com>
OpenClaw-Publication: 58c87505-0064-4391-994e-57587279bcca

* fix(voice): use explicit Node message transfer lists

Replace six Node postMessage lint suppressions with explicit empty transfer lists. Preserve the Node messaging contract and the unchanged production suppression allowlist. Reproduced the original CI shard plan before the fix; the same plan, focused worker messaging tests, and type-aware lint pass afterward.

Co-authored-by: steipete <58493+steipete@users.noreply.github.com>

* fix(openai): keep worker runtime out of cold voice catalogs

Load the default media-socket factory only through abort-aware connection admission. Preserve injected factories and WebRTC routing, and centralize unchanged request-ID construction in the wire owner. Keep the cold-catalog contract and line cap intact.

Co-authored-by: steipete <58493+steipete@users.noreply.github.com>

* refactor(openai): keep socket contracts below connection admission

Move unchanged socket interfaces into the existing shared leaf and migrate all private consumers, preserving lazy worker loading without a type import cycle.

Co-authored-by: steipete <58493+steipete@users.noreply.github.com>

* fix(discord): preserve direct speech ownership through playback

Track exact speech with an epoch-scoped shared playback witness. Order completion behind bounded direct-port flush receipts, including no-audio responses, and prevent stale idle or flush events from retiring newer speech. Cover original failures and the pending-flush race with deterministic production-owner regressions.

Co-authored-by: steipete <58493+steipete@users.noreply.github.com>

* fix(discord): preserve audio ownership across worker retirement

Worked on by:
- @steipete

Co-authored-by: steipete <58493+steipete@users.noreply.github.com>
OpenClaw-Publication: 6294ffa7-0c2f-456b-9baf-3dcdab60efc8

* refactor(discord): keep the worker status slot private

Worked on by:
- @steipete

Co-authored-by: steipete <58493+steipete@users.noreply.github.com>
OpenClaw-Publication: dde5cff1-21c6-4bfa-825e-47c2f1ec2393

---------

Co-authored-by: steipete <58493+steipete@users.noreply.github.com>
2026-09-20 21:49:35 -07:00

171 lines
5.5 KiB
TypeScript

import {
createStreamingPcmResampler,
mulawToPcm,
pcmToMulaw,
type RealtimeVoiceBridgeCreateRequest,
} from "openclaw/plugin-sdk/realtime-voice-provider";
const RELAY_FRAME_SAMPLES = 480;
const MAX_PENDING_RELAY_FRAMES = 250;
export const OPENAI_QUICKSILVER_AUDIO_FRAME_DURATION_MS = 20;
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 assertOpenAIQuicksilverPcmOutput(
format: RealtimeVoiceBridgeCreateRequest["audioFormat"],
): void {
if (
format &&
(format.encoding !== "pcm16" || format.sampleRateHz !== 24_000 || format.channels !== 1)
) {
throw new Error("GPT-Live direct audio output requires mono PCM16 at 24 kHz");
}
}
/** Keeps telephony resampling state and its delayed output tail with the audio adapter. */
export class OpenAIQuicksilverAudioAdapter {
private readonly telephony: boolean;
private inbound = createStreamingPcmResampler(8_000, 24_000);
private outbound = createStreamingPcmResampler(24_000, 8_000);
constructor(
private readonly config: Pick<RealtimeVoiceBridgeCreateRequest, "audioFormat" | "onAudio">,
) {
this.telephony = config.audioFormat?.encoding === "g711_ulaw";
}
decodeInput(audio: Buffer): Buffer {
return this.telephony ? this.inbound.process(mulawToPcm(audio)) : audio;
}
sendOutput(pcm: Buffer): void {
const audio = this.telephony ? pcmToMulaw(this.outbound.process(pcm)) : pcm;
if (audio.length > 0) {
this.config.onAudio(audio);
}
}
finishOutput(): void {
if (!this.telephony) {
return;
}
const tail = pcmToMulaw(this.outbound.flush());
this.outbound = createStreamingPcmResampler(24_000, 8_000);
if (tail.length > 0) {
this.config.onAudio(tail);
}
}
reset(): void {
this.inbound = createStreamingPcmResampler(8_000, 24_000);
this.outbound = createStreamingPcmResampler(24_000, 8_000);
}
}
/** One real-time clock for WebSocket PCM and WebRTC RTP; stalls never drain capture in a burst. */
export class OpenAIQuicksilverAudioClock {
private timer: ReturnType<typeof setTimeout> | undefined;
constructor(private readonly onFrame: (skippedFrames: number) => void) {}
start(): void {
if (this.timer) {
return;
}
let nextFrameAt = performance.now();
const tick = () => {
const skippedFrames = Math.max(
0,
Math.floor((performance.now() - nextFrameAt) / OPENAI_QUICKSILVER_AUDIO_FRAME_DURATION_MS),
);
nextFrameAt += (skippedFrames + 1) * OPENAI_QUICKSILVER_AUDIO_FRAME_DURATION_MS;
// Publish before delivery so synchronous teardown also cancels the first tick's successor.
this.timer = setTimeout(tick, Math.max(0, nextFrameAt - performance.now()));
this.timer.unref?.();
this.onFrame(skippedFrames);
};
tick();
}
stop(): void {
if (this.timer) {
clearTimeout(this.timer);
this.timer = undefined;
}
}
}
export class OpenAIQuicksilverPendingAudio {
private storage: Buffer | undefined;
private readOffset = 0;
private pendingBytes = 0;
get length(): number {
return this.pendingBytes;
}
append(incoming: Buffer): void {
const evenLength = incoming.length - (incoming.length % 2);
if (evenLength === 0) {
return;
}
// Capture owns its input only until the next callback; copy each retained sample
// once into the circular tail instead of copying the entire history per frame.
const retainedBytes = Math.min(evenLength, OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES);
const sourceOffset = evenLength - retainedBytes;
const storage = (this.storage ??= Buffer.alloc(OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES));
const droppedBytes = Math.max(
0,
this.pendingBytes + retainedBytes - OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES,
);
this.readOffset = (this.readOffset + droppedBytes) % OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES;
this.pendingBytes -= droppedBytes;
const writeOffset =
(this.readOffset + this.pendingBytes) % OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES;
const firstBytes = Math.min(
retainedBytes,
OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES - writeOffset,
);
incoming.copy(storage, writeOffset, sourceOffset, sourceOffset + firstBytes);
if (firstBytes < retainedBytes) {
incoming.copy(storage, 0, sourceOffset + firstBytes, sourceOffset + retainedBytes);
}
this.pendingBytes += retainedBytes;
}
readInto(target: Buffer): number {
const evenLength = target.length - (target.length % 2);
const readBytes = Math.min(evenLength, this.pendingBytes);
const storage = this.storage;
if (readBytes === 0 || !storage) {
return 0;
}
const firstBytes = Math.min(
readBytes,
OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES - this.readOffset,
);
storage.copy(target, 0, this.readOffset, this.readOffset + firstBytes);
if (firstBytes < readBytes) {
storage.copy(target, firstBytes, 0, readBytes - firstBytes);
}
this.readOffset = (this.readOffset + readBytes) % OPENAI_QUICKSILVER_MAX_PENDING_AUDIO_BYTES;
this.pendingBytes -= readBytes;
if (this.pendingBytes === 0) {
this.readOffset = 0;
}
return readBytes;
}
clear(): void {
this.storage = undefined;
this.readOffset = 0;
this.pendingBytes = 0;
}
}