Skip to content

Commit 5ceb1bd

Browse files
authored
Fix hx-sse reader-cancel race that produces an unhandled rejection (#3935)
Fix reader-cancel race in hx-sse.js parseSSE's `finally` released the reader's lock, but `connection.reader` was only nulled one microtask later (via the for-await-of loop's async-generator close protocol). Any concurrent `connection.reader.cancel()` call (pauseOnBackground's visibilitychange handler, or cleanup()) landing in that one-microtask window canceled an already-released reader, which per the Streams spec rejects instead of throwing synchronously -- surfacing as an unhandled promise rejection. Fix: pass `connection` into parseSSE and null `connection.reader` in the same synchronous finally block as releaseLock(), closing the window instead of guarding each call site. Includes a regression test that deterministically lands in the race window by chaining the concurrent cancel trigger off the same promise the internal `await reader.read()` is waiting on.
1 parent 243d4b1 commit 5ceb1bd

2 files changed

Lines changed: 53 additions & 3 deletions

File tree

src/ext/hx-sse.js

Lines changed: 7 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -32,7 +32,9 @@
3232
// SSE PARSER
3333
// ========================================
3434

35-
async function* parseSSE(reader, lastEventId = '') {
35+
async function* parseSSE(connection) {
36+
let reader = connection.reader;
37+
let lastEventId = connection.lastEventId;
3638
let decoder = new TextDecoder();
3739
let buffer = '';
3840
let hasData = false;
@@ -100,6 +102,9 @@
100102
}
101103
} finally {
102104
reader.releaseLock();
105+
// On the same tick as releaseLock(), so cleanup() or visibilityHandler
106+
// calling connection.reader?.cancel() can't hit a released reader.
107+
connection.reader = null;
103108
}
104109
}
105110

@@ -241,7 +246,7 @@
241246
try {
242247
connection.reader = currentResponse.body.getReader();
243248

244-
for await (let msg of parseSSE(connection.reader, connection.lastEventId)) {
249+
for await (let msg of parseSSE(connection)) {
245250
if (!element.isConnected || reconnectRequested) break;
246251

247252
if (msg.hasId) {
@@ -285,7 +290,6 @@
285290
}
286291
}
287292

288-
connection.reader = null;
289293
if (!element.isConnected) break;
290294

291295
connection.attempt++;

test/tests/ext/hx-sse.js

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -591,6 +591,52 @@ describe('hx-sse SSE extension', function() {
591591
stream.close();
592592
});
593593

594+
it('cancelling an already-released reader does not produce an unhandled rejection', async function() {
595+
// Hand-rolled reader so the test can fire visibilitychange (→ cancel()) off the same
596+
// promise the extension's `await reader.read()` is waiting on, landing it deterministically
597+
// in the same tick the stream closes on.
598+
let resolveRead;
599+
let released = false;
600+
const pendingRead = new Promise(resolve => { resolveRead = resolve; });
601+
const reader = {
602+
read: () => pendingRead,
603+
releaseLock() { released = true; },
604+
cancel() {
605+
if (released) return Promise.reject(new TypeError('Cannot cancel a stream using a released reader'));
606+
released = true;
607+
return Promise.resolve();
608+
}
609+
};
610+
611+
fetchMock.mockResponse('GET', '/reader-cancel-race', () => {
612+
const response = new MockResponse({getReader: () => reader}, {
613+
headers: {'Content-Type': 'text/event-stream'}
614+
});
615+
response.body = {getReader: () => reader};
616+
return response;
617+
});
618+
619+
createProcessedHTML('<button hx-get="/reader-cancel-race" hx-config="sse.pauseOnBackground:true" hx-swap="innerHTML">Connect</button>');
620+
621+
let rejections = [];
622+
let onRejection = e => rejections.push(e.reason);
623+
window.addEventListener('unhandledrejection', onRejection);
624+
625+
find('button').click();
626+
await waitForEvent('htmx:after:sse:connection');
627+
628+
pendingRead.then(() => {
629+
Object.defineProperty(document, 'hidden', {value: true, configurable: true});
630+
document.dispatchEvent(new Event('visibilitychange'));
631+
});
632+
resolveRead({done: true, value: undefined});
633+
634+
await htmx.timeout(20);
635+
window.removeEventListener('unhandledrejection', onRejection);
636+
637+
assert.deepEqual(rejections, [], 'Attempting to cancel an already-released reader must not produce an unhandled rejection');
638+
});
639+
594640
it('custom events trigger on element and bubble', async function() {
595641
const stream = mockStreamResponse('/custom-events');
596642
createProcessedHTML('<button hx-get="/custom-events" hx-swap="innerHTML">Connect</button>');

0 commit comments

Comments
 (0)