From 54d682e6f5f44487922de2526c3ee4aa9fab7321 Mon Sep 17 00:00:00 2001 From: Yusuke Wada Date: Tue, 15 Jul 2025 17:22:50 +0900 Subject: [PATCH 1/2] fix: handle disconnection gracefully without canceling stream --- src/utils.ts | 48 +++++++++++++++++++++++++++++++++------------- test/utils.test.ts | 41 ++++++++++++++++++++++++++++++++++++++- 2 files changed, 75 insertions(+), 14 deletions(-) diff --git a/src/utils.ts b/src/utils.ts index c704c93..c6839bc 100644 --- a/src/utils.ts +++ b/src/utils.ts @@ -5,38 +5,60 @@ export function writeFromReadableStream(stream: ReadableStream, writ if (stream.locked) { throw new TypeError('ReadableStream is locked.') } else if (writable.destroyed) { - stream.cancel() return } + const reader = stream.getReader() - writable.on('close', cancel) - writable.on('error', cancel) - reader.read().then(flow, cancel) + let clientDisconnected = false + + const handleClientDisconnect = () => { + clientDisconnected = true + } + + const handleError = () => { + clientDisconnected = true + } + + writable.on('close', handleClientDisconnect) + writable.on('error', handleError) + + reader.read().then(flow, handleStreamError) + return reader.closed.finally(() => { - writable.off('close', cancel) - writable.off('error', cancel) + writable.off('close', handleClientDisconnect) + writable.off('error', handleError) }) + // eslint-disable-next-line @typescript-eslint/no-explicit-any - function cancel(error?: any) { - reader.cancel(error).catch(() => {}) - if (error) { - writable.destroy(error) + function handleStreamError(error: any) { + if (!clientDisconnected) { + if (error) { + writable.destroy(error) + } } } + function onDrain() { - reader.read().then(flow, cancel) + if (!clientDisconnected) { + reader.read().then(flow, handleStreamError) + } } + function flow({ done, value }: ReadableStreamReadResult): void | Promise { + if (clientDisconnected) { + return + } + try { if (done) { writable.end() } else if (!writable.write(value)) { writable.once('drain', onDrain) } else { - return reader.read().then(flow, cancel) + return reader.read().then(flow, handleStreamError) } } catch (e) { - cancel(e) + handleStreamError(e) } } } diff --git a/test/utils.test.ts b/test/utils.test.ts index 3673c5f..53f6972 100644 --- a/test/utils.test.ts +++ b/test/utils.test.ts @@ -1,4 +1,5 @@ -import { buildOutgoingHttpHeaders } from '../src/utils' +import { Writable } from 'node:stream' +import { buildOutgoingHttpHeaders, writeFromReadableStream } from '../src/utils' describe('buildOutgoingHttpHeaders', () => { it('original content-type is preserved', () => { @@ -71,3 +72,41 @@ describe('buildOutgoingHttpHeaders', () => { }) }) }) + +describe('writeFromReadableStream', () => { + it('should handle client disconnection gracefully without canceling stream', async () => { + let enqueueCalled = false + let cancelCalled = false + + // Create test ReadableStream + const stream = new ReadableStream({ + start(controller) { + setTimeout(() => { + try { + controller.enqueue(new TextEncoder().encode('test')) + enqueueCalled = true + } catch { + // Test should fail if error occurs + } + controller.close() + }, 100) + }, + cancel() { + cancelCalled = true + }, + }) + + // Test Writable stream + const writable = new Writable() + + // Simulate client disconnection after 50ms + setTimeout(() => { + writable.destroy() + }, 50) + + await writeFromReadableStream(stream, writable) + + expect(enqueueCalled).toBe(true) // enqueue should succeed + expect(cancelCalled).toBe(false) // cancel should not be called + }) +}) From e90d5963f4a1d282928beb54e83fdfe8d80d17b1 Mon Sep 17 00:00:00 2001 From: Yusuke Wada Date: Sat, 19 Jul 2025 19:19:29 +0900 Subject: [PATCH 2/2] Simplify Co-authored-by: Taku Amano --- src/utils.ts | 23 ++++------------------- 1 file changed, 4 insertions(+), 19 deletions(-) diff --git a/src/utils.ts b/src/utils.ts index c6839bc..1f4240d 100644 --- a/src/utils.ts +++ b/src/utils.ts @@ -9,46 +9,31 @@ export function writeFromReadableStream(stream: ReadableStream, writ } const reader = stream.getReader() - let clientDisconnected = false - - const handleClientDisconnect = () => { - clientDisconnected = true - } const handleError = () => { - clientDisconnected = true + // ignore the error } - writable.on('close', handleClientDisconnect) writable.on('error', handleError) reader.read().then(flow, handleStreamError) return reader.closed.finally(() => { - writable.off('close', handleClientDisconnect) writable.off('error', handleError) }) // eslint-disable-next-line @typescript-eslint/no-explicit-any function handleStreamError(error: any) { - if (!clientDisconnected) { - if (error) { - writable.destroy(error) - } + if (error) { + writable.destroy(error) } } function onDrain() { - if (!clientDisconnected) { - reader.read().then(flow, handleStreamError) - } + reader.read().then(flow, handleStreamError) } function flow({ done, value }: ReadableStreamReadResult): void | Promise { - if (clientDisconnected) { - return - } - try { if (done) { writable.end()