diff --git a/CHANGELOG.md b/CHANGELOG.md index 534142f5a..8ccf8a9ac 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -14,6 +14,8 @@ metadata and the backend fallback mirror it. ### Fixed +- Release Stories and Audiobook response streams when a render event fails or a pending read is stopped — thanks @rudycelekli! (#2652) + - MCP speech tools wait through model loading and progress-extended CPU renders instead of timing out before the backend (#2609) ## [0.5.7] — 2026-10-05 diff --git a/docs/expressive-speech.md b/docs/expressive-speech.md index 7666f56ac..33b466635 100644 --- a/docs/expressive-speech.md +++ b/docs/expressive-speech.md @@ -258,3 +258,5 @@ promised date — when each lands, this page gets updated in the same PR. Omitted API join options keep legacy hard joins (zero gaps, no edge trim). Clients enable seamless joins by sending their chosen gaps and trim explicitly. + +Stories and Audiobook stream consumers cancel and unlock their response reader when an event handler fails or rendering stops. A supplied abort signal also wakes a pending read, so the backend can observe the disconnect and stop scheduling chapters. diff --git a/electron/src/shared/test/longformStream.test.js b/electron/src/shared/test/longformStream.test.js index c60a5a24a..78502c280 100644 --- a/electron/src/shared/test/longformStream.test.js +++ b/electron/src/shared/test/longformStream.test.js @@ -62,6 +62,42 @@ describe('consumeLongformStream', () => { expect(events.length).toBe(1); }); + it('cancels and unlocks the native stream when an event handler rejects a render', async () => { + let cancelled = false; + const body = new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode(sse({ type: 'error', error: 'render failed' }))); + }, + cancel() { cancelled = true; }, + }); + const failure = new Error('render failed'); + await expect(consumeLongformStream(new Response(body), () => { throw failure; })) + .rejects.toBe(failure); + expect(cancelled).toBe(true); + expect(body.locked).toBe(false); + }); + + it('unlocks a native stream after normal completion', async () => { + const body = new ReadableStream({ start(controller) { controller.close(); } }); + await consumeLongformStream(new Response(body), () => {}); + expect(body.locked).toBe(false); + }); + + it('wakes a pending native read when the supplied signal aborts', async () => { + const ctrl = new AbortController(); + let cancelled = false; + const body = new ReadableStream({ cancel() { cancelled = true; } }); + const promise = consumeLongformStream(new Response(body), () => {}, { signal: ctrl.signal }); + await Promise.resolve(); + ctrl.abort(); + await expect(Promise.race([ + promise, + new Promise((_, reject) => setTimeout(() => reject(new Error('read stayed pending')), 50)), + ])).resolves.toBeUndefined(); + expect(cancelled).toBe(true); + expect(body.locked).toBe(false); + }); + it('throws when the response has no body', async () => { await expect(consumeLongformStream({}, () => {})).rejects.toThrow(/no response stream/); }); diff --git a/electron/src/shared/utils/longformStream.js b/electron/src/shared/utils/longformStream.js index c1fe98646..3183d70c4 100644 --- a/electron/src/shared/utils/longformStream.js +++ b/electron/src/shared/utils/longformStream.js @@ -42,26 +42,27 @@ export async function consumeLongformStream(res, onEvent, { isAborted, signal } // Releasing the reader cancels the underlying stream (closes the fetch), so // the server sees the disconnect. Best-effort: a stream already closed/errored // — or a caller's fake reader without cancel() — must not throw here. - const releaseStream = async () => { + let cancellation; + const releaseStream = () => (cancellation ??= (async () => { try { await reader.cancel(); } catch { /* already closed/errored, or no cancel() — nothing to release */ } - }; + })()); + const onAbort = () => { void releaseStream(); }; + signal?.addEventListener('abort', onAbort, { once: true }); try { while (true) { - if (aborted()) { - await releaseStream(); - return; - } + if (aborted()) return; const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); const { lines, rest } = splitSSEBuffer(buffer); buffer = rest; for (const line of lines) { + if (aborted()) return; const evt = parseSSELine(line); if (evt) onEvent(evt); } @@ -70,10 +71,11 @@ export async function consumeLongformStream(res, onEvent, { isAborted, signal } // An abort mid-read (AbortController.abort() / reader.cancel()) rejects the // pending read() — swallow it when WE initiated the stop; re-throw a genuine // stream/transport error so callers still surface it. - if (aborted()) { - await releaseStream(); - return; - } + if (aborted()) return; throw e; + } finally { + signal?.removeEventListener('abort', onAbort); + await releaseStream(); + reader.releaseLock?.(); } }