Skip to content

Commit 0f31254

Browse files
committed
Merge pull request openclaw#116696 from openclaw/fix/talk-gateway-uplink-bound
fix(ui): bound Talk relay microphone uplink * refs/remotes/origin/pr/116696: fix(ui): bound Talk relay microphone uplink
2 parents fcf2e7a + a84ede7 commit 0f31254

2 files changed

Lines changed: 152 additions & 8 deletions

File tree

ui/src/pages/chat/realtime-talk-gateway-relay.test.ts

Lines changed: 110 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -380,6 +380,116 @@ describe("GatewayRelayRealtimeTalkTransport", () => {
380380
expect(onInputLevel).toHaveBeenLastCalledWith(0);
381381
});
382382

383+
it("bounds stalled microphone appends and aborts every owner on stop", async () => {
384+
const onStatus = vi.fn();
385+
const client = createClient();
386+
let activeAppends = 0;
387+
let peakActiveAppends = 0;
388+
const appendSignals: AbortSignal[] = [];
389+
vi.mocked(client["request"]).mockImplementation((method, _params, options) => {
390+
if (method !== "talk.session.appendAudio") {
391+
return Promise.resolve({});
392+
}
393+
const signal = options?.signal;
394+
if (!signal) {
395+
return Promise.reject(new Error("missing append abort signal"));
396+
}
397+
appendSignals.push(signal);
398+
activeAppends += 1;
399+
peakActiveAppends = Math.max(peakActiveAppends, activeAppends);
400+
return new Promise((_, reject) => {
401+
signal.addEventListener(
402+
"abort",
403+
() => {
404+
activeAppends -= 1;
405+
reject(new Error("append aborted"));
406+
},
407+
{ once: true },
408+
);
409+
});
410+
});
411+
const transport = createTransport({ callbacks: { onStatus }, client });
412+
413+
await transport.start();
414+
const samples = new Float32Array(4096);
415+
for (let index = 0; index < 10_000; index += 1) {
416+
pumpMicrophone(samples);
417+
}
418+
419+
const appendCalls = requestCallsFor(client, "talk.session.appendAudio");
420+
expect(appendCalls).toHaveLength(4);
421+
expect(peakActiveAppends).toBe(4);
422+
expect(activeAppends).toBe(4);
423+
expect(new Set(appendSignals).size).toBe(1);
424+
expect(
425+
appendCalls.every(
426+
(call) => call[2]?.signal === appendSignals[0] && call[2]?.timeoutMs === 8_000,
427+
),
428+
).toBe(true);
429+
430+
transport.stop();
431+
transport.stop();
432+
await Promise.resolve();
433+
434+
expect(activeAppends).toBe(0);
435+
expect(appendSignals.every((signal) => signal.aborted)).toBe(true);
436+
expect(requestCallsFor(client, "talk.session.close")).toHaveLength(1);
437+
expect(onStatus).not.toHaveBeenCalled();
438+
});
439+
440+
it("preserves accepted microphone frame order", async () => {
441+
const client = createClient();
442+
const transport = createTransport({ client });
443+
444+
await transport.start();
445+
for (const timestamp of [10, 20, 30, 40]) {
446+
audioCurrentTime = timestamp / 1_000;
447+
pumpMicrophone(new Float32Array(4096));
448+
}
449+
450+
expect(
451+
requestCallsFor(client, "talk.session.appendAudio").map(
452+
(call) => (call[1] as { timestamp: number }).timestamp,
453+
),
454+
).toEqual([10, 20, 30, 40]);
455+
transport.stop();
456+
});
457+
458+
it("ignores a stale append rejection after a replacement starts", async () => {
459+
const oldStatus = vi.fn();
460+
const oldClient = createClient();
461+
let rejectOldAppend: (error: Error) => void = () => undefined;
462+
vi.mocked(oldClient["request"]).mockImplementation((method) => {
463+
if (method !== "talk.session.appendAudio") {
464+
return Promise.resolve({});
465+
}
466+
return new Promise((_, reject) => {
467+
rejectOldAppend = reject;
468+
});
469+
});
470+
const oldTransport = createTransport({ callbacks: { onStatus: oldStatus }, client: oldClient });
471+
472+
await oldTransport.start();
473+
pumpMicrophone(new Float32Array(4096));
474+
oldTransport.stop();
475+
476+
const replacementStatus = vi.fn();
477+
const replacementClient = createClient();
478+
const replacement = createTransport({
479+
callbacks: { onStatus: replacementStatus },
480+
client: replacementClient,
481+
});
482+
await replacement.start();
483+
pumpMicrophone(new Float32Array(4096));
484+
rejectOldAppend(new Error("late stale append failure"));
485+
await Promise.resolve();
486+
487+
expect(requestCallsFor(replacementClient, "talk.session.appendAudio")).toHaveLength(1);
488+
expect(oldStatus).not.toHaveBeenCalled();
489+
expect(replacementStatus).not.toHaveBeenCalled();
490+
replacement.stop();
491+
});
492+
383493
it("stops microphone pumping when the relay rejects appended audio", async () => {
384494
const onStatus = vi.fn();
385495
const client = createClient();

ui/src/pages/chat/realtime-talk-gateway-relay.ts

Lines changed: 42 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@ import {
2222
const BARGE_IN_RMS_THRESHOLD = 0.02;
2323
const BARGE_IN_PEAK_THRESHOLD = 0.08;
2424
const BARGE_IN_CONSECUTIVE_SPEECH_FRAMES = 2;
25+
const MAX_PENDING_AUDIO_APPENDS = 4;
26+
const AUDIO_APPEND_TIMEOUT_MS = 8_000;
2527

2628
export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport {
2729
private media: MediaStream | null = null;
@@ -31,6 +33,8 @@ export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport
3133
private readonly inputPump = new RealtimeTalkPcmInputPump();
3234
private unsubscribe: (() => void) | null = null;
3335
private closed = false;
36+
private audioAppendAbortController: AbortController | null = null;
37+
private readonly pendingAudioAppends = new Set<Promise<unknown>>();
3438
private readonly outputQueue = new RealtimeTalkPcmOutputQueue();
3539
private readonly consultAbortControllers = new Map<string, AbortController>();
3640
private readonly completedToolCalls = new Set<string>();
@@ -80,6 +84,8 @@ export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport
8084
this.media = media;
8185
this.inputContext = new AudioContext({ sampleRate: this.session.audio.inputSampleRateHz });
8286
this.outputContext = new AudioContext({ sampleRate: this.session.audio.outputSampleRateHz });
87+
this.abortPendingAudioAppends();
88+
this.audioAppendAbortController = new AbortController();
8389
if (this.ctx.callbacks.onInputLevel) {
8490
this.inputMeter = new RealtimeTalkMediaStreamMeter(this.ctx.callbacks.onInputLevel);
8591
this.inputMeter.start(this.media, this.inputContext);
@@ -104,6 +110,7 @@ export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport
104110
this.unsubscribe?.();
105111
this.unsubscribe = null;
106112
this.inputPump.stop();
113+
this.abortPendingAudioAppends();
107114
this.inputMeter?.stop();
108115
this.inputMeter = null;
109116
// Mark callbacks recurse until playback drains, so shutdown must cancel every owned timer.
@@ -128,28 +135,55 @@ export class GatewayRelayRealtimeTalkTransport implements RealtimeTalkTransport
128135
if (this.closed) {
129136
return;
130137
}
131-
const pcm = floatToPcm16(samples);
132138
if (this.detectBargeInSpeech(samples)) {
133139
this.cancelOutputForBargeIn();
134140
}
135-
void this.ctx.client
136-
.request("talk.session.appendAudio", {
137-
sessionId: this.session.relaySessionId,
138-
audioBase64: bytesToBase64(pcm),
139-
timestamp: Math.round((this.inputContext?.currentTime ?? 0) * 1000),
140-
})
141+
const abortController = this.audioAppendAbortController;
142+
// Live microphone frames become stale once the Gateway falls behind, so drop new
143+
// frames at the ownership cap instead of growing a latency queue.
144+
if (
145+
!abortController ||
146+
abortController.signal.aborted ||
147+
this.pendingAudioAppends.size >= MAX_PENDING_AUDIO_APPENDS
148+
) {
149+
return;
150+
}
151+
const pcm = floatToPcm16(samples);
152+
const request = this.ctx.client
153+
.request(
154+
"talk.session.appendAudio",
155+
{
156+
sessionId: this.session.relaySessionId,
157+
audioBase64: bytesToBase64(pcm),
158+
timestamp: Math.round((this.inputContext?.currentTime ?? 0) * 1000),
159+
},
160+
{
161+
signal: abortController.signal,
162+
timeoutMs: AUDIO_APPEND_TIMEOUT_MS,
163+
},
164+
)
141165
.catch((error: unknown) => {
142-
if (!this.closed) {
166+
if (!this.closed && !abortController.signal.aborted) {
143167
this.ctx.callbacks.onStatus?.(
144168
"error",
145169
error instanceof Error ? error.message : String(error),
146170
);
147171
this.stop();
148172
}
149173
});
174+
this.pendingAudioAppends.add(request);
175+
void request.finally(() => {
176+
this.pendingAudioAppends.delete(request);
177+
});
150178
});
151179
}
152180

181+
private abortPendingAudioAppends(): void {
182+
this.audioAppendAbortController?.abort();
183+
this.audioAppendAbortController = null;
184+
this.pendingAudioAppends.clear();
185+
}
186+
153187
private handleRelayEvent(event: GatewayRelayEvent): void {
154188
if (event.relaySessionId !== this.session.relaySessionId || this.closed) {
155189
return;

0 commit comments

Comments
 (0)