Skip to content

Commit f5419b5

Browse files
committed
fix(openrouter): release music stream readers
1 parent 14fd10f commit f5419b5

2 files changed

Lines changed: 123 additions & 41 deletions

File tree

extensions/openrouter/music-generation-provider.test.ts

Lines changed: 88 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -39,19 +39,69 @@ vi.mock("openclaw/plugin-sdk/provider-http", async (importOriginal) => {
3939
};
4040
});
4141

42-
function sseResponse(lines: string[]): Response {
42+
function sseResponse(lines: string[], options?: { releaseLock?: () => void }): Response {
4343
const encoder = new TextEncoder();
44-
return new Response(
45-
new ReadableStream({
46-
start(controller) {
47-
for (const line of lines) {
48-
controller.enqueue(encoder.encode(line));
49-
}
50-
controller.close();
51-
},
52-
}),
53-
{ status: 200, headers: { "content-type": "text/event-stream" } },
54-
);
44+
if (!options?.releaseLock) {
45+
return new Response(
46+
new ReadableStream({
47+
start(controller) {
48+
for (const line of lines) {
49+
controller.enqueue(encoder.encode(line));
50+
}
51+
controller.close();
52+
},
53+
}),
54+
{ status: 200, headers: { "content-type": "text/event-stream" } },
55+
);
56+
}
57+
58+
const chunks: Array<ReadableStreamReadResult<Uint8Array>> = lines.map((line) => ({
59+
done: false,
60+
value: encoder.encode(line),
61+
}));
62+
chunks.push({ done: true, value: undefined });
63+
const reader = {
64+
read: async () => chunks.shift() ?? { done: true, value: undefined },
65+
cancel: async () => undefined,
66+
releaseLock: options.releaseLock,
67+
} as ReadableStreamDefaultReader<Uint8Array>;
68+
69+
return {
70+
ok: true,
71+
status: 200,
72+
headers: new Headers({ "content-type": "text/event-stream" }),
73+
body: {
74+
getReader: () => reader,
75+
},
76+
} as Response;
77+
}
78+
79+
function sseResponseLines(params: {
80+
audio?: string;
81+
transcript?: string;
82+
done?: boolean;
83+
}): string[] {
84+
const lines: string[] = [];
85+
if (params.audio || params.transcript) {
86+
lines.push(
87+
`data: ${JSON.stringify({
88+
choices: [
89+
{
90+
delta: {
91+
audio: {
92+
...(params.audio ? { data: params.audio } : {}),
93+
...(params.transcript ? { transcript: params.transcript } : {}),
94+
},
95+
},
96+
},
97+
],
98+
})}\n`,
99+
);
100+
}
101+
if (params.done) {
102+
lines.push("data: [DONE]\n");
103+
}
104+
return lines;
55105
}
56106

57107
function stalledSseResponse(line: string): Response {
@@ -150,6 +200,32 @@ describe("openrouter music generation provider", () => {
150200
expect(release).toHaveBeenCalledOnce();
151201
});
152202

203+
it("releases OpenRouter audio stream readers after completion", async () => {
204+
const releaseLock = vi.fn();
205+
postJsonRequestMock.mockResolvedValue({
206+
response: sseResponse(
207+
sseResponseLines({
208+
audio: Buffer.from("wav-bytes").toString("base64"),
209+
done: true,
210+
}),
211+
{ releaseLock },
212+
),
213+
release: vi.fn(async () => {}),
214+
});
215+
216+
await expect(
217+
buildOpenRouterMusicGenerationProvider().generateMusic({
218+
provider: "openrouter",
219+
model: "google/lyria-3-pro-preview",
220+
prompt: "release stream reader",
221+
cfg: {},
222+
}),
223+
).resolves.toMatchObject({
224+
tracks: [{ mimeType: "audio/wav" }],
225+
});
226+
expect(releaseLock).toHaveBeenCalledTimes(1);
227+
});
228+
153229
it("decodes independently padded OpenRouter audio chunks", async () => {
154230
postJsonRequestMock.mockResolvedValue({
155231
response: sseResponse([

extensions/openrouter/music-generation-provider.ts

Lines changed: 35 additions & 29 deletions
Original file line numberDiff line numberDiff line change
@@ -179,40 +179,46 @@ async function readOpenRouterAudioStream(
179179
const result = { audioBuffers: [] as Buffer[], transcriptChunks: [] as string[] };
180180
let buffer = "";
181181
let doneSeen = false;
182-
for (;;) {
183-
const { value, done } = await readOpenRouterStreamChunk(reader, deadline);
184-
if (done) {
185-
break;
186-
}
187-
buffer += decoder.decode(value, { stream: true });
188-
const lines = buffer.split(/\r?\n/u);
189-
buffer = lines.pop() ?? "";
190-
for (const line of lines) {
191-
if (processOpenRouterSseLine(line.trim(), result)) {
192-
await reader.cancel();
193-
return {
194-
audioBuffer: Buffer.concat(result.audioBuffers),
195-
transcript: result.transcriptChunks.join(""),
196-
};
182+
try {
183+
for (;;) {
184+
const { value, done } = await readOpenRouterStreamChunk(reader, deadline);
185+
if (done) {
186+
break;
187+
}
188+
buffer += decoder.decode(value, { stream: true });
189+
const lines = buffer.split(/\r?\n/u);
190+
buffer = lines.pop() ?? "";
191+
for (const line of lines) {
192+
if (processOpenRouterSseLine(line.trim(), result)) {
193+
await reader.cancel();
194+
return {
195+
audioBuffer: Buffer.concat(result.audioBuffers),
196+
transcript: result.transcriptChunks.join(""),
197+
};
198+
}
197199
}
198200
}
199-
}
200-
resolveOpenRouterStreamRemainingMs(deadline);
201-
buffer += decoder.decode();
202-
if (buffer.trim()) {
203-
for (const line of buffer.split(/\r?\n/u)) {
204-
if (processOpenRouterSseLine(line.trim(), result)) {
205-
doneSeen = true;
201+
resolveOpenRouterStreamRemainingMs(deadline);
202+
buffer += decoder.decode();
203+
if (buffer.trim()) {
204+
for (const line of buffer.split(/\r?\n/u)) {
205+
if (processOpenRouterSseLine(line.trim(), result)) {
206+
doneSeen = true;
207+
}
206208
}
207209
}
210+
if (!doneSeen) {
211+
throw new Error("OpenRouter music generation stream ended before completion");
212+
}
213+
return {
214+
audioBuffer: Buffer.concat(result.audioBuffers),
215+
transcript: result.transcriptChunks.join(""),
216+
};
217+
} finally {
218+
try {
219+
reader.releaseLock();
220+
} catch {}
208221
}
209-
if (!doneSeen) {
210-
throw new Error("OpenRouter music generation stream ended before completion");
211-
}
212-
return {
213-
audioBuffer: Buffer.concat(result.audioBuffers),
214-
transcript: result.transcriptChunks.join(""),
215-
};
216222
}
217223

218224
export function buildOpenRouterMusicGenerationProvider(): MusicGenerationProvider {

0 commit comments

Comments
 (0)