diff --git a/.changeset/two-kids-cry.md b/.changeset/two-kids-cry.md new file mode 100644 index 0000000..9d267b0 --- /dev/null +++ b/.changeset/two-kids-cry.md @@ -0,0 +1,5 @@ +--- +'fetch-nodeshim': patch +--- + +Propagate errors for duplex request/response streams, and ensure early errors propagate to the Response stream diff --git a/src/__tests__/fetch-errors.test.ts b/src/__tests__/fetch-errors.test.ts new file mode 100644 index 0000000..f27615c --- /dev/null +++ b/src/__tests__/fetch-errors.test.ts @@ -0,0 +1,107 @@ +import { describe, it, expect, beforeEach, afterEach } from 'vitest'; +import { Readable } from 'node:stream'; + +import TestServer from './utils/server.js'; +import { fetch } from '../fetch'; + +describe('fetch error handling', () => { + const server = new TestServer(); + + beforeEach(() => server.start()); + afterEach(() => server.stop()); + + describe('incoming stream errors', () => { + it('should propagate error when connection resets mid-body', async () => { + const url = server.mock((_req, res) => { + res.writeHead(200, { 'Content-Length': '1000' }); + res.write('partial'); + setTimeout(() => res.destroy(), 10); + }); + + const response = await fetch(url); + expect(response.ok).toBe(true); + await expect(response.text()).rejects.toThrow(); + }); + + it('should not cause unhandled errors on incoming stream', async () => { + const errors: Error[] = []; + const handler = (e: Error) => errors.push(e); + process.on('uncaughtException', handler); + + const url = server.mock((_req, res) => { + res.writeHead(200, { 'Transfer-Encoding': 'chunked' }); + res.write('data'); + setTimeout(() => res.destroy(), 10); + }); + + const response = await fetch(url); + try { + await response.text(); + } catch {} + await new Promise(r => setTimeout(r, 50)); + + process.off('uncaughtException', handler); + expect(errors).toHaveLength(0); + }); + }); + + describe('request body pipeline errors', () => { + it('should reject when request body stream errors', async () => { + const body = new Readable({ + read() { + this.push('data'); + setTimeout(() => this.destroy(new Error('stream error')), 10); + }, + }); + + const url = server.mock((req, res) => { + req.on('data', () => {}); + req.on('end', () => res.end('ok')); + req.on('error', () => {}); + }); + + await expect(fetch(url, { method: 'POST', body })).rejects.toThrow( + 'stream error' + ); + }); + }); + + describe('decompression errors', () => { + it('should propagate gzip decompression errors', async () => { + const url = server.mock((_req, res) => { + res.writeHead(200, { 'Content-Encoding': 'gzip' }); + res.end('not gzip'); + }); + + const response = await fetch(url); + expect(response.ok).toBe(true); + await expect(response.text()).rejects.toThrow(); + }); + + it('should propagate brotli decompression errors', async () => { + const url = server.mock((_req, res) => { + res.writeHead(200, { 'Content-Encoding': 'br' }); + res.end('not brotli'); + }); + + const response = await fetch(url); + await expect(response.text()).rejects.toThrow(); + }); + }); + + describe('abort handling', () => { + it('should abort during response streaming', async () => { + const controller = new AbortController(); + + const url = server.mock((_req, res) => { + res.writeHead(200, { 'Transfer-Encoding': 'chunked' }); + const id = setInterval(() => res.write('x'), 10); + res.on('close', () => clearInterval(id)); + }); + + const response = await fetch(url, { signal: controller.signal }); + setTimeout(() => controller.abort(), 30); + await expect(response.text()).rejects.toThrow(); + }); + }); +}); diff --git a/src/__tests__/fetch.test.ts b/src/__tests__/fetch.test.ts index d3b5db5..9797b9b 100644 --- a/src/__tests__/fetch.test.ts +++ b/src/__tests__/fetch.test.ts @@ -595,7 +595,7 @@ describe(fetch, () => { }); setTimeout(() => controller.abort(), 100); await expect(response$).rejects.toThrowErrorMatchingInlineSnapshot( - `[AbortError: The operation was aborted]` + `[AbortError: This operation was aborted]` ); }); @@ -613,10 +613,10 @@ describe(fetch, () => { ]; setTimeout(() => controller.abort(), 100); await expect(fetches[0]).rejects.toThrowErrorMatchingInlineSnapshot( - `[AbortError: The operation was aborted]` + `[AbortError: This operation was aborted]` ); await expect(fetches[1]).rejects.toThrowErrorMatchingInlineSnapshot( - `[AbortError: The operation was aborted]` + `[AbortError: This operation was aborted]` ); }); @@ -627,7 +627,7 @@ describe(fetch, () => { signal: controller.signal, }); }).rejects.toThrowErrorMatchingInlineSnapshot( - `[AbortError: The operation was aborted]` + `[AbortError: This operation was aborted]` ); }); @@ -639,7 +639,7 @@ describe(fetch, () => { await expect(() => fetch(request) ).rejects.toThrowErrorMatchingInlineSnapshot( - `[AbortError: The operation was aborted]` + `[AbortError: This operation was aborted]` ); }); @@ -672,7 +672,7 @@ describe(fetch, () => { }); controller.abort(); await expect(response$).rejects.toThrowErrorMatchingInlineSnapshot( - `[AbortError: The operation was aborted]` + `[AbortError: This operation was aborted]` ); }); @@ -708,7 +708,7 @@ describe(fetch, () => { controller.abort(); await bodyError$; await expect(response$).rejects.toMatchInlineSnapshot( - `[AbortError: The operation was aborted]` + `[AbortError: This operation was aborted]` ); }); diff --git a/src/__tests__/utils/server.js b/src/__tests__/utils/server.js index 25cfe32..8ea628c 100644 --- a/src/__tests__/utils/server.js +++ b/src/__tests__/utils/server.js @@ -3,7 +3,7 @@ import http from 'http'; import zlib from 'zlib'; import Busboy from 'busboy'; -import {once} from 'events'; +import { once } from 'events'; export default class TestServer { constructor() { @@ -42,15 +42,20 @@ export default class TestServer { return `http://${this.hostname}:${this.port}/mocked`; } + mock(handler) { + this.server.nextResponseHandler = handler; + return `http://${this.hostname}:${this.port}/mocked`; + } + router(request, res) { const p = request.url; if (p === '/mocked') { if (this.nextResponseHandler) { - this.nextResponseHandler(res); + this.nextResponseHandler(request, res); this.nextResponseHandler = undefined; } else { - throw new Error('No mocked response. Use ’TestServer.mockResponse()’.'); + throw new Error("No mocked response. Use 'TestServer.mockResponse()'."); } } @@ -91,9 +96,11 @@ export default class TestServer { if (p === '/json') { res.statusCode = 200; res.setHeader('Content-Type', 'application/json'); - res.end(JSON.stringify({ - name: 'value' - })); + res.end( + JSON.stringify({ + name: 'value', + }) + ); } if (p === '/gzip') { @@ -244,8 +251,8 @@ export default class TestServer { } if (p === '/redirect/301/rn') { - res.statusCode = 301 - res.setHeader('Location', '/403') + res.statusCode = 301; + res.setHeader('Location', '/403'); res.write('301 Permanently moved.\r\n'); res.end(); } @@ -321,7 +328,7 @@ export default class TestServer { if (p === '/redirect/chunked') { res.writeHead(301, { Location: '/inspect', - 'Transfer-Encoding': 'chunked' + 'Transfer-Encoding': 'chunked', }); setTimeout(() => res.end(), 10); } @@ -349,7 +356,7 @@ export default class TestServer { } if (p === '/error/premature') { - res.writeHead(200, {'content-length': 50}); + res.writeHead(200, { 'content-length': 50 }); res.write('foo'); setTimeout(() => { res.destroy(); @@ -359,13 +366,13 @@ export default class TestServer { if (p === '/error/premature/chunked') { res.writeHead(200, { 'Content-Type': 'application/json', - 'Transfer-Encoding': 'chunked' + 'Transfer-Encoding': 'chunked', }); - res.write(`${JSON.stringify({data: 'hi'})}\n`); + res.write(`${JSON.stringify({ data: 'hi' })}\n`); setTimeout(() => { - res.write(`${JSON.stringify({data: 'bye'})}\n`); + res.write(`${JSON.stringify({ data: 'bye' })}\n`); }, 50); setTimeout(() => { @@ -435,38 +442,43 @@ export default class TestServer { body += c; }); request.on('end', () => { - res.end(JSON.stringify({ - inspect: true, - method: request.method, - url: request.url, - headers: request.headers, - body - })); + res.end( + JSON.stringify({ + inspect: true, + method: request.method, + url: request.url, + headers: request.headers, + body, + }) + ); }); } if (p === '/multipart') { res.statusCode = 200; res.setHeader('Content-Type', 'application/json'); - const busboy = new Busboy({headers: request.headers}); + const busboy = new Busboy({ headers: request.headers }); let body = ''; busboy.on('file', async (fieldName, file, fileName) => { body += `${fieldName}=${fileName}`; // consume file data // eslint-disable-next-line no-empty, no-unused-vars - for await (const c of file) { } + for await (const c of file) { + } }); busboy.on('field', (fieldName, value) => { body += `${fieldName}=${value}`; }); busboy.on('finish', () => { - res.end(JSON.stringify({ - method: request.method, - url: request.url, - headers: request.headers, - body - })); + res.end( + JSON.stringify({ + method: request.method, + url: request.url, + headers: request.headers, + body, + }) + ); }); request.pipe(busboy); } diff --git a/src/fetch.ts b/src/fetch.ts index c8f766f..cd539cf 100644 --- a/src/fetch.ts +++ b/src/fetch.ts @@ -147,9 +147,27 @@ async function _fetch( const protocol = requestOptions.protocol === 'https:' ? https : http; const outgoing = protocol.request(requestOptions); - outgoing.on('response', incoming => { + let incoming: http.IncomingMessage | undefined; + + const destroy = (reason?: any) => { + if (reason) { + outgoing?.destroy(signal?.aborted ? signal.reason : reason); + incoming?.destroy(signal?.aborted ? signal.reason : reason); + reject(signal?.aborted ? signal.reason : reason); + } + }; + + signal?.addEventListener('abort', destroy); + + outgoing.on('response', _incoming => { + if (signal?.aborted) { + return; + } + + incoming = _incoming; incoming.setTimeout(0); // Forcefully disable timeout incoming.socket.unref(); + incoming.on('error', destroy); const init = { status: incoming.statusCode, @@ -209,16 +227,6 @@ async function _fetch( } } - const destroy = (reason?: any) => { - signal?.removeEventListener('abort', destroy); - if (reason) { - incoming.destroy(signal?.aborted ? signal.reason : reason); - reject(signal?.aborted ? signal.reason : reason); - } - }; - - signal?.addEventListener('abort', destroy); - let body: Readable | null = incoming; const encoding = init.headers.get('Content-Encoding')?.toLowerCase(); if (method === 'HEAD' || init.status === 204 || init.status === 304) { @@ -226,6 +234,7 @@ async function _fetch( } else if (encoding != null) { init.headers.set('Content-Encoding', encoding); body = pipeline(body, createContentDecoder(encoding), destroy); + outgoing.on('error', destroy); } resolve( @@ -237,7 +246,7 @@ async function _fetch( ); }); - outgoing.on('error', reject); + outgoing.on('error', destroy); if (!requestHeaders.has('Accept')) requestHeaders.set('Accept', '*/*'); if (requestBody.contentType) @@ -261,9 +270,7 @@ async function _fetch( requestBody.body instanceof Stream ? requestBody.body : Readable.fromWeb(requestBody.body); - pipeline(body, outgoing, error => { - if (error) reject(error); - }); + pipeline(body, outgoing, destroy); } }