From 443554bae93cd044232fbbf886359a6174a375e9 Mon Sep 17 00:00:00 2001 From: Vincent Koc Date: Fri, 31 Jul 2026 07:23:35 +0800 Subject: [PATCH] fix(meeting-bot): bound realtime egress --- src/meeting-bot/realtime-engine.test.ts | 86 ++++++++--- src/meeting-bot/realtime-engine.ts | 184 +++++++++++++++++++++--- 2 files changed, 233 insertions(+), 37 deletions(-) diff --git a/src/meeting-bot/realtime-engine.test.ts b/src/meeting-bot/realtime-engine.test.ts index cd06bd1be1ec..62e948d014c5 100644 --- a/src/meeting-bot/realtime-engine.test.ts +++ b/src/meeting-bot/realtime-engine.test.ts @@ -133,6 +133,20 @@ describe("meeting realtime engine output ownership", () => { } }); + it("does not start an invalidated write after a same-turn clear", async () => { + const fixture = await createEngineFixture(); + try { + fixture.callbacks.onAudio(Buffer.from([1, 2, 3])); + fixture.callbacks.onClearAudio(); + await Promise.resolve(); + + expect(fixture.writeOutput).not.toHaveBeenCalled(); + expect(fixture.clearOutput).toHaveBeenCalled(); + } finally { + await fixture.handle.stop(); + } + }); + it.each(["response.done", "response.cancelled"])( "bounds queued bytes and rejects stale output through %s", async (terminalType) => { @@ -145,17 +159,19 @@ describe("meeting realtime engine output ownership", () => { const fresh = Buffer.from([5]); fixture.callbacks.onAudio(first); - fixture.callbacks.onAudio(queued); - fixture.callbacks.onAudio(overflow); - - expect(fixture.handleBargeIn).toHaveBeenCalledWith({ - audioPlaybackActive: true, - force: true, - }); - expect(fixture.clearOutput).toHaveBeenCalledOnce(); await vi.waitFor(() => { expect(fixture.writeOutput).toHaveBeenCalledTimes(1); }); + fixture.callbacks.onAudio(queued); + fixture.callbacks.onAudio(overflow); + + await vi.waitFor(() => { + expect(fixture.handleBargeIn).toHaveBeenCalledWith({ + audioPlaybackActive: true, + force: true, + }); + }); + expect(fixture.clearOutput).toHaveBeenCalledOnce(); expect(fixture.writeOutput).toHaveBeenLastCalledWith(first); fixture.callbacks.onAudio(late); @@ -167,6 +183,10 @@ describe("meeting realtime engine output ownership", () => { expect(fixture.writeOutput).toHaveBeenCalledTimes(2); }); expect(fixture.writeOutput).toHaveBeenLastCalledWith(fresh); + expect(fixture.clearOutput).toHaveBeenCalledTimes(2); + expect(fixture.clearOutput.mock.invocationCallOrder[1]).toBeLessThan( + fixture.writeOutput.mock.invocationCallOrder[1] ?? 0, + ); expect(fixture.beginOutput).toHaveBeenCalledTimes(2); fixture.releaseWrite(1); } finally { @@ -178,18 +198,50 @@ describe("meeting realtime engine output ownership", () => { it("bounds queued tiny-frame ownership", async () => { const fixture = await createEngineFixture(); try { - for (let index = 0; index < 257; index += 1) { - fixture.callbacks.onAudio(Buffer.from([index])); - } - - expect(fixture.handleBargeIn).toHaveBeenCalledWith({ - audioPlaybackActive: true, - force: true, - }); - expect(fixture.clearOutput).toHaveBeenCalledOnce(); + fixture.callbacks.onAudio(Buffer.from([0])); await vi.waitFor(() => { expect(fixture.writeOutput).toHaveBeenCalledTimes(1); }); + for (let index = 1; index < 257; index += 1) { + fixture.callbacks.onAudio(Buffer.from([index])); + } + + await vi.waitFor(() => { + expect(fixture.handleBargeIn).toHaveBeenCalledWith({ + audioPlaybackActive: true, + force: true, + }); + }); + expect(fixture.clearOutput).toHaveBeenCalledOnce(); + fixture.releaseWrite(0); + } finally { + await fixture.handle.stop(); + } + }); + + it("does not report a deferred cancellation race after response completion", async () => { + const fixture = await createEngineFixture(); + try { + fixture.callbacks.onAudio(Buffer.alloc(48_000, 1)); + await vi.waitFor(() => { + expect(fixture.writeOutput).toHaveBeenCalledTimes(1); + }); + fixture.callbacks.onAudio(Buffer.alloc(48_000, 2)); + fixture.callbacks.onAudio(Buffer.from([3])); + await vi.waitFor(() => { + expect(fixture.handleBargeIn).toHaveBeenCalledOnce(); + }); + + fixture.callbacks.onEvent?.({ direction: "server", type: "response.done" }); + fixture.callbacks.onEvent?.({ + direction: "server", + type: "error", + detail: "Cancellation failed: no active response found", + }); + + expect(fixture.handle.getHealth().recentTalkEvents.map((event) => event.type)).not.toContain( + "session.error", + ); fixture.releaseWrite(0); } finally { await fixture.handle.stop(); diff --git a/src/meeting-bot/realtime-engine.ts b/src/meeting-bot/realtime-engine.ts index adfee1df799b..2bf65d5f55f0 100644 --- a/src/meeting-bot/realtime-engine.ts +++ b/src/meeting-bot/realtime-engine.ts @@ -101,6 +101,9 @@ export const MEETING_AGENT_TRANSCRIPT_DEBOUNCE_MS = 900; // Playback duration plus a tail blocks live loopback; transcript lookback catches delayed echo. export const MEETING_OUTPUT_ECHO_SUPPRESSION_TAIL_MS = 3_000; export const MEETING_TRANSCRIPT_ECHO_LOOKBACK_MS = 45_000; +const MEETING_REALTIME_OUTPUT_MAX_PENDING_MS = 2_000; +const MEETING_REALTIME_OUTPUT_MAX_PENDING_FRAMES = 256; +const MEETING_REALTIME_CANCELLATION_RACE_DETAIL = "Cancellation failed: no active response found"; export function meetingOutputBytesPerMs(audioFormat: MeetingRealtimeAudioFormat): number { return audioFormat === "g711-ulaw-8khz" ? 8 : 48; @@ -301,10 +304,36 @@ export async function startMeetingRealtimeEngine(params: { let realtimeReady = false; let lastClearAt: string | undefined; let clearCount = 0; + let outputGeneration = 0; + let outputWriteActive = false; + let outputClearPending = 0; + let outputClearAfterActive = false; + let outputPendingBytes = 0; + let outputPendingFrames = 0; + let outputBlocked: { token: symbol } | undefined; + let outputGenerationActive = false; + let outputClearTail = Promise.resolve(); + const outputQueue: Array<{ audio: Buffer; generation: number }> = []; + const outputMaxPendingBytes = + meetingOutputBytesPerMs(params.config.chrome.audioFormat) * + MEETING_REALTIME_OUTPUT_MAX_PENDING_MS; const realtimeLogScope = params.logPrefix ? `${params.logPrefix} realtime` : "realtime"; + const invalidateOutputQueue = () => { + outputGeneration += 1; + outputQueue.length = 0; + outputPendingBytes = 0; + outputPendingFrames = 0; + outputGenerationActive = false; + }; + const stop = async () => { - stopped = true; + if (!stopped) { + stopped = true; + outputBlocked = undefined; + outputClearAfterActive = false; + invalidateOutputQueue(); + } if (stopPromise) { await stopPromise; return; @@ -360,28 +389,126 @@ export async function startMeetingRealtimeEngine(params: { ); }); }; - const clearOutputPlayback = () => { + const clearOutputPlayback = (): Promise => { if (stopped) { - return; + return Promise.resolve(); } clearCount += 1; lastClearAt = new Date().toISOString(); - void params.transport.clearOutput().catch((error: unknown) => { - params.logger.warn( - `${params.platform.logScope} ${params.logPrefix ? `${params.logPrefix} audio clear` : "audio output clear"} failed: ${formatErrorMessage(error)}`, - ); - stopAfterFailure("audio output clear"); + outputClearPending += 1; + const clear = outputClearTail + .then(async () => { + if (!stopped) { + await params.transport.clearOutput(); + } + }) + .catch((error: unknown) => { + params.logger.warn( + `${params.platform.logScope} ${params.logPrefix ? `${params.logPrefix} audio clear` : "audio output clear"} failed: ${formatErrorMessage(error)}`, + ); + stopAfterFailure("audio output clear"); + }) + .finally(() => { + outputClearPending -= 1; + pumpOutputQueue(); + }); + outputClearTail = clear; + return clear; + }; + + const pumpOutputQueue = () => { + if (stopped || outputWriteActive || outputClearPending > 0) { + return; + } + const next = outputQueue.shift(); + if (!next) { + return; + } + if (next.generation !== outputGeneration) { + pumpOutputQueue(); + return; + } + outputWriteActive = true; + void Promise.resolve() + .then(() => { + if (stopped || next.generation !== outputGeneration) { + return; + } + return params.transport.writeOutput(next.audio); + }) + .catch((error: unknown) => { + if (stopped || next.generation !== outputGeneration) { + return; + } + params.logger.warn( + `${params.platform.logScope} ${params.logPrefix ? `${params.logPrefix} audio output` : "audio output"} failed: ${formatErrorMessage(error)}`, + ); + stopAfterFailure("audio output"); + }) + .finally(() => { + outputWriteActive = false; + if (next.generation === outputGeneration) { + outputPendingBytes -= next.audio.byteLength; + outputPendingFrames -= 1; + } + if (outputClearAfterActive && !stopped) { + outputClearAfterActive = false; + void clearOutputPlayback(); + return; + } + pumpOutputQueue(); + }); + }; + + const blockOutput = (): { blocked: boolean; token: symbol } => { + if (outputBlocked) { + return { blocked: false, token: outputBlocked.token }; + } + const token = Symbol("meeting-realtime-output-blocked"); + outputBlocked = { token }; + outputClearAfterActive = outputWriteActive; + invalidateOutputQueue(); + return { blocked: true, token }; + }; + + const handleOutputBackpressure = () => { + const pendingBytes = outputPendingBytes; + const pendingFrames = outputPendingFrames; + const { blocked, token } = blockOutput(); + if (!blocked) { + return; + } + params.logger.warn( + `${params.platform.logScope} ${realtimeLogScope} audio output backpressured: pendingBytes=${pendingBytes} pendingFrames=${pendingFrames}`, + ); + harness.flushOutput(clearOutputPlayback); + harness.finishOutputAudio("output-backpressure"); + queueMicrotask(() => { + if (stopped || outputBlocked?.token !== token) { + return; + } + harness.handleBargeIn({ audioPlaybackActive: true, force: true }, () => {}); }); }; - const writeOutputAudio = (audio: Buffer) => { - void params.transport.writeOutput(audio).catch((error: unknown) => { - params.logger.warn( - `${params.platform.logScope} ${params.logPrefix ? `${params.logPrefix} audio output` : "audio output"} failed: ${formatErrorMessage(error)}`, - ); - stopAfterFailure("audio output"); - }); + + const queueOutputAudio = (audio: Buffer): boolean => { + if (stopped || outputBlocked) { + return false; + } + if ( + audio.byteLength > outputMaxPendingBytes - outputPendingBytes || + outputPendingFrames >= MEETING_REALTIME_OUTPUT_MAX_PENDING_FRAMES + ) { + handleOutputBackpressure(); + return false; + } + outputPendingBytes += audio.byteLength; + outputPendingFrames += 1; + outputQueue.push({ audio, generation: outputGeneration }); + pumpOutputQueue(); + return true; }; - let outputGenerationActive = false; + const startHumanBargeInMonitor = () => { if (!params.transport.startBargeInMonitor) { return; @@ -499,18 +626,21 @@ export async function startMeetingRealtimeEngine(params: { audioSink: { isOpen: () => !stopped, sendAudio: (audio) => { + if (!queueOutputAudio(audio)) { + return; + } if (!outputGenerationActive) { params.transport.beginOutput?.(); outputGenerationActive = true; } harness.outputActivity.markPlaybackStarted(); harness.recordOutputAudio(audio); - writeOutputAudio(audio); }, clearAudio: () => { - outputGenerationActive = false; - harness.flushOutput(clearOutputPlayback); - harness.finishOutputAudio("clear"); + if (blockOutput().blocked) { + harness.flushOutput(clearOutputPlayback); + harness.finishOutputAudio("clear"); + } }, }, onTranscript: (role, text, isFinal) => { @@ -576,9 +706,23 @@ export async function startMeetingRealtimeEngine(params: { final: true, }); } else if (event.type === "response.done") { + outputBlocked = undefined; outputGenerationActive = false; harness.finishOutputAudio("response.done"); harness.endTurn("response.done"); + } else if (outputBlocked && event.type === "response.cancelled") { + outputBlocked = undefined; + outputGenerationActive = false; + harness.finishOutputAudio(event.type); + } else if ( + event.type === "error" && + event.detail === MEETING_REALTIME_CANCELLATION_RACE_DETAIL + ) { + if (outputBlocked) { + outputBlocked = undefined; + outputGenerationActive = false; + harness.finishOutputAudio(event.type); + } } else if (event.type === "error") { harness.emit({ type: "session.error",