Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 2 additions & 0 deletions docs/expressive-speech.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
36 changes: 36 additions & 0 deletions electron/src/shared/test/longformStream.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -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/);
});
Expand Down
22 changes: 12 additions & 10 deletions electron/src/shared/utils/longformStream.js
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand All @@ -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?.();
}
}
Loading