fix(elevenlabs): stop oversized streaming TTS audio (#109621)

* fix(elevenlabs): cap streamed TTS audio size

* fix(elevenlabs): avoid stream cap cancellation race

* fix(elevenlabs): release bounded stream reader locks

* fix(elevenlabs): release stream resources on early cleanup

* fix(elevenlabs): satisfy typed stream cleanup lint
This commit is contained in:
xingzhou
2026-07-18 05:15:08 +01:00
committed by GitHub
parent 3659c85e53
commit c6abd2e4d4
2 changed files with 198 additions and 7 deletions
+106 -5
View File
@@ -1,4 +1,5 @@
// Elevenlabs tests cover tts plugin behavior.
import { MAX_AUDIO_BYTES } from "openclaw/plugin-sdk/media-runtime";
import { afterEach, describe, expect, it, vi } from "vitest";
import { createStreamingErrorResponse } from "../test-support/streaming-error-response.js";
import { elevenLabsTTS, elevenLabsTTSStream } from "./tts.js";
@@ -203,12 +204,112 @@ describe("elevenlabs tts diagnostics", () => {
...createDefaultTtsRequest(),
latencyTier: 2,
});
try {
const url = getUrlFromFirstFetchCall(fetchMock);
expect(url.pathname).toBe("/v1/text-to-speech/pMsXgVXv3BLzUgSXRplE/stream");
expect(url.searchParams.get("optimize_streaming_latency")).toBe("2");
const reader = result.audioStream.getReader();
await expect(reader.read()).resolves.toEqual({
done: false,
value: new Uint8Array([1, 2, 3]),
});
await expect(reader.read()).resolves.toEqual({ done: true, value: undefined });
expect(audioStream.locked).toBe(false);
} finally {
await result.release();
}
});
const url = getUrlFromFirstFetchCall(fetchMock);
expect(url.pathname).toBe("/v1/text-to-speech/pMsXgVXv3BLzUgSXRplE/stream");
expect(url.searchParams.get("optimize_streaming_latency")).toBe("2");
expect(result.audioStream).toBeInstanceOf(ReadableStream);
await result.release();
it("releases an unread provider stream when the public release hook runs", async () => {
const cancel = vi.fn();
const audioStream = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new Uint8Array([1, 2, 3]));
},
cancel,
});
const fetchMock = vi.fn(
async () => new Response(audioStream, { headers: { "content-type": "audio/mpeg" } }),
);
globalThis.fetch = fetchMock as unknown as typeof fetch;
const result = await elevenLabsTTSStream(createDefaultTtsRequest());
try {
expect(audioStream.locked).toBe(true);
await result.release();
await result.release();
expect(cancel).toHaveBeenCalledOnce();
expect(audioStream.locked).toBe(false);
} finally {
await result.audioStream.cancel().catch(() => undefined);
await result.release();
}
});
it("releases a partially consumed provider stream while its consumer holds a reader", async () => {
const cancel = vi.fn();
const audioStream = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new Uint8Array([1, 2, 3]));
},
cancel,
});
const fetchMock = vi.fn(
async () => new Response(audioStream, { headers: { "content-type": "audio/mpeg" } }),
);
globalThis.fetch = fetchMock as unknown as typeof fetch;
const result = await elevenLabsTTSStream(createDefaultTtsRequest());
const reader = result.audioStream.getReader();
try {
await expect(reader.read()).resolves.toEqual({
done: false,
value: new Uint8Array([1, 2, 3]),
});
await result.release();
expect(cancel).toHaveBeenCalledOnce();
expect(audioStream.locked).toBe(false);
await expect(reader.read()).resolves.toEqual({ done: true, value: undefined });
} finally {
await reader.cancel().catch(() => undefined);
reader.releaseLock();
await result.release();
}
});
it("cancels streamed audio before delivering bytes beyond the audio limit", async () => {
const cancel = vi.fn();
const audioStream = new ReadableStream<Uint8Array>({
start(controller) {
controller.enqueue(new Uint8Array(MAX_AUDIO_BYTES));
controller.enqueue(new Uint8Array([1]));
},
cancel,
});
const fetchMock = vi.fn(
async () => new Response(audioStream, { headers: { "content-type": "audio/mpeg" } }),
);
globalThis.fetch = fetchMock as unknown as typeof fetch;
const result = await elevenLabsTTSStream(createDefaultTtsRequest());
try {
const reader = result.audioStream.getReader();
const first = await reader.read();
expect(first.done).toBe(false);
expect(first.value).toHaveLength(MAX_AUDIO_BYTES);
await expect(reader.read()).rejects.toThrow(
`ElevenLabs API error: audio response exceeds ${MAX_AUDIO_BYTES} bytes`,
);
expect(cancel).toHaveBeenCalledOnce();
expect(audioStream.locked).toBe(false);
} finally {
await result.release();
}
});
it("rejects JSON success stream responses as malformed audio", async () => {
+92 -2
View File
@@ -1,4 +1,5 @@
// Elevenlabs plugin module implements tts behavior.
import { MAX_AUDIO_BYTES } from "openclaw/plugin-sdk/media-runtime";
import {
assertOkOrThrowProviderError,
assertProviderBinaryResponseContent,
@@ -48,6 +49,83 @@ function normalizeElevenLabsLatencyTier(latencyTier: number | undefined): number
return latencyTier;
}
// Mirror the buffered cap without buffering. Own the reader because Node can leak
// transform writer rejections when playback cancellation races an overflow.
function createBoundedElevenLabsAudioStream(stream: ReadableStream<Uint8Array>): {
audioStream: ReadableStream<Uint8Array>;
release: () => Promise<void>;
} {
let reader: ReadableStreamDefaultReader<Uint8Array> | undefined;
let totalBytes = 0;
const releaseReader = (activeReader: ReadableStreamDefaultReader<Uint8Array>) => {
if (reader !== activeReader) {
return;
}
reader = undefined;
activeReader.releaseLock();
};
const cancelReader = async (reason?: unknown) => {
const activeReader = reader;
if (!activeReader) {
return;
}
try {
await activeReader.cancel(reason).catch(() => undefined);
} finally {
releaseReader(activeReader);
}
};
const audioStream = new ReadableStream<Uint8Array>({
start() {
reader = stream.getReader();
},
async pull(controller) {
const activeReader = reader;
if (!activeReader) {
controller.close();
return;
}
try {
const chunk = await activeReader.read();
if (chunk.done) {
releaseReader(activeReader);
controller.close();
return;
}
const remainingBytes = MAX_AUDIO_BYTES - totalBytes;
if (chunk.value.byteLength > remainingBytes) {
if (remainingBytes > 0) {
controller.enqueue(chunk.value.subarray(0, remainingBytes));
}
const error = new Error(
`ElevenLabs API error: audio response exceeds ${MAX_AUDIO_BYTES} bytes`,
);
await activeReader.cancel(error).catch(() => undefined);
releaseReader(activeReader);
controller.error(error);
return;
}
totalBytes += chunk.value.byteLength;
controller.enqueue(chunk.value);
} catch (error) {
releaseReader(activeReader);
controller.error(error);
}
},
async cancel(reason) {
await cancelReader(reason);
},
});
return {
audioStream,
release: () => cancelReader(new Error("ElevenLabs TTS stream released")),
};
}
type ElevenLabsTtsRequestParams = {
text: string;
apiKey: string;
@@ -191,10 +269,22 @@ export async function elevenLabsTTSStream(params: ElevenLabsTtsRequestParams): P
if (!response.body) {
throw new Error("ElevenLabs API response missing audio stream");
}
const boundedStream = createBoundedElevenLabsAudioStream(response.body);
let releasePromise: Promise<void> | undefined;
const releaseAll = () => {
releasePromise ??= (async () => {
try {
await boundedStream.release();
} finally {
await release();
}
})();
return releasePromise;
};
handedOff = true;
return {
audioStream: response.body,
release,
audioStream: boundedStream.audioStream,
release: releaseAll,
};
} finally {
if (!handedOff) {