|
| 1 | +/** |
| 2 | + * @vitest-environment node |
| 3 | + */ |
| 4 | +import { afterAll, afterEach, beforeAll, expect, it } from 'vitest' |
| 5 | +import http from 'node:http' |
| 6 | +import net, { AddressInfo } from 'node:net' |
| 7 | +import { ClientRequestInterceptor } from '../../../../src/interceptors/ClientRequest' |
| 8 | + |
| 9 | +const largeResponseBody = Buffer.alloc(1024 * 1024, 'a') |
| 10 | +const interceptor = new ClientRequestInterceptor() |
| 11 | + |
| 12 | +const rawHttpServer = net.createServer((socket) => { |
| 13 | + socket.once('data', () => { |
| 14 | + socket.write('HTTP/1.1 200 OK\r\n') |
| 15 | + socket.write('Connection: close\r\n') |
| 16 | + // Use conventional response framing. The regression is caused by the |
| 17 | + // large response body and delayed consumption, not by this header. |
| 18 | + socket.write(`Content-Length: ${largeResponseBody.byteLength}\r\n`) |
| 19 | + socket.write('Content-Type: text/plain\r\n') |
| 20 | + socket.write('\r\n') |
| 21 | + socket.end(largeResponseBody) |
| 22 | + }) |
| 23 | +}) |
| 24 | + |
| 25 | +beforeAll(async () => { |
| 26 | + interceptor.apply() |
| 27 | + await new Promise<void>((resolve) => { |
| 28 | + rawHttpServer.listen(0, '127.0.0.1', resolve) |
| 29 | + }) |
| 30 | +}) |
| 31 | + |
| 32 | +afterEach(() => { |
| 33 | + interceptor.removeAllListeners() |
| 34 | +}) |
| 35 | + |
| 36 | +afterAll(async () => { |
| 37 | + interceptor.dispose() |
| 38 | + await new Promise<void>((resolve, reject) => { |
| 39 | + rawHttpServer.close((error) => { |
| 40 | + if (error) { |
| 41 | + reject(error) |
| 42 | + return |
| 43 | + } |
| 44 | + |
| 45 | + resolve() |
| 46 | + }) |
| 47 | + }) |
| 48 | +}) |
| 49 | + |
| 50 | +it('delivers a large passthrough response to a delayed consumer before closing', async () => { |
| 51 | + expect.assertions(2) |
| 52 | + |
| 53 | + const address = rawHttpServer.address() as AddressInfo |
| 54 | + const result = await new Promise<{ |
| 55 | + receivedBodySize: number |
| 56 | + isComplete: boolean |
| 57 | + }>((resolve, reject) => { |
| 58 | + const request = http.get( |
| 59 | + `http://127.0.0.1:${address.port}/resource`, |
| 60 | + (response) => { |
| 61 | + let bytesRead = 0 |
| 62 | + |
| 63 | + response.on('error', reject) |
| 64 | + |
| 65 | + // Delay response consumption by one tick so the original socket close |
| 66 | + // can race with this passthrough socket's readable completion. |
| 67 | + response.pause() |
| 68 | + |
| 69 | + setImmediate(() => { |
| 70 | + response.on('data', (chunk: Buffer) => { |
| 71 | + bytesRead += chunk.byteLength |
| 72 | + }) |
| 73 | + response.on('end', () => { |
| 74 | + resolve({ |
| 75 | + receivedBodySize: bytesRead, |
| 76 | + isComplete: response.complete, |
| 77 | + }) |
| 78 | + }) |
| 79 | + response.resume() |
| 80 | + }) |
| 81 | + } |
| 82 | + ) |
| 83 | + |
| 84 | + request.on('error', reject) |
| 85 | + }) |
| 86 | + |
| 87 | + expect(result.receivedBodySize).toBe(largeResponseBody.byteLength) |
| 88 | + expect(result.isComplete).toBe(true) |
| 89 | +}) |
0 commit comments