From 9b2d92661f8a025cfaa267913f8e39c8ddfa84bb Mon Sep 17 00:00:00 2001 From: PekingSpades <180665176+PekingSpades@users.noreply.github.com> Date: Tue, 15 Sep 2026 20:40:34 +0800 Subject: [PATCH] fix: drain unobserved response body streams Consume response tracking streams when no complete response listener needs buffering, preventing active responses from growing an unused queue outside the configured body size limit. --- src/server/mockttp-server.ts | 3 +- src/util/request-utils.ts | 7 ++ test/response-buffering.spec.ts | 118 ++++++++++++++++++++++++++++++++ 3 files changed, 127 insertions(+), 1 deletion(-) create mode 100644 test/response-buffering.spec.ts diff --git a/src/server/mockttp-server.ts b/src/server/mockttp-server.ts index 693d82212..ac18e527e 100644 --- a/src/server/mockttp-server.ts +++ b/src/server/mockttp-server.ts @@ -717,12 +717,14 @@ export class MockttpServer extends AbstractMockttp implements Mockttp { this.announceInitialRequestAsync(emitter, request); + const hasResponseListener = emitter.listenerCount('response') > 0; const response = trackResponse( rawResponse, request.timingEvents, request.tags, { maxSize: this.maxBodySize, + captureBody: hasResponseListener, onWriteHead: () => this.announceInitialResponseAsync(emitter, response), onBodyData: emitter.listenerCount('response-body-data') > 0 ? this.announceBodyDataAsync.bind(this, emitter, 'response') @@ -738,7 +740,6 @@ export class MockttpServer extends AbstractMockttp implements Mockttp { response.sendInformationalResponse(100, []); } - const hasResponseListener = emitter.listenerCount('response') > 0; if (hasResponseListener) { // Start buffering response body if there's somebody who // might want to hear about it later diff --git a/src/util/request-utils.ts b/src/util/request-utils.ts index ba3b9b328..2ebc12b89 100644 --- a/src/util/request-utils.ts +++ b/src/util/request-utils.ts @@ -542,6 +542,7 @@ export function trackResponse( tags: string[], options: { maxSize: number, + captureBody?: boolean, onWriteHead: () => void, onBodyData?: (id: string, timestamp: number, content: Uint8Array, isEnded: boolean) => void, onInformationalResponse: (status: number, flatHeaders: string[]) => void @@ -591,6 +592,12 @@ export function trackResponse( emitBodyDataEvents(trackedResponse, trackingStream, options.onBodyData); } + if (options.captureBody === false) { + // maxSize only limits an active body reader. Without one, the side + // stream must still be consumed to avoid an unbounded writable queue. + trackingStream.resume(); + } + const originalWriteHeader = trackedResponse.writeHead; const originalWrite = trackedResponse.write; const originalEnd = trackedResponse.end; diff --git a/test/response-buffering.spec.ts b/test/response-buffering.spec.ts new file mode 100644 index 000000000..9b9779387 --- /dev/null +++ b/test/response-buffering.spec.ts @@ -0,0 +1,118 @@ +import { Buffer } from 'buffer'; +import * as http from 'http'; +import * as stream from 'stream'; + +import { trackResponse } from '../src/util/request-utils'; +import { expect, nodeOnly, getDeferred } from './test-utils'; + +nodeOnly(() => { + describe("Response body buffering", () => { + let trackingStreams: stream.PassThrough[]; + let responses: stream.Writable[]; + + beforeEach(() => { + trackingStreams = []; + responses = []; + }); + afterEach(() => { + responses.forEach(response => response.destroy()); + trackingStreams.forEach(tracking => tracking.destroy()); + }); + + function createResponse( + captureBody: boolean, + maxSize = 1024, + onBodyData?: Parameters[3]['onBodyData'] + ) { + let writtenBytes = 0; + const sink = Object.assign(new stream.Writable({ + write(chunk, _encoding, callback) { + writtenBytes += chunk.length; + callback(); + } + }), { + getHeaders: () => ({}), + writeHead() {}, + addTrailers() {} + }); + responses.push(sink); + + // Capture the side stream only while constructing the response tracker. + // Queue lengths give a deterministic bound without relying on GC or RSS. + const descriptor = Object.getOwnPropertyDescriptor(stream, 'PassThrough')!; + const OriginalPassThrough = stream.PassThrough; + Object.defineProperty(stream, 'PassThrough', { + ...descriptor, + value: class extends OriginalPassThrough { + constructor() { + super(); + trackingStreams.push(this); + } + } + }); + const options = { + captureBody, + maxSize, + onBodyData, + onWriteHead() {}, + onInformationalResponse() {} + }; + try { + const response = trackResponse( + sink as unknown as http.ServerResponse, + { startTime: Date.now(), startTimestamp: performance.now() }, + [], + options + ); + return { response, tracking: trackingStreams[trackingStreams.length - 1], bytes: () => writtenBytes }; + } finally { + Object.defineProperty(stream, 'PassThrough', descriptor); + } + } + + it("does not queue unobserved streaming responses beyond maxBodySize", async () => { + const { response, tracking, bytes } = createResponse(false); + for (let i = 0; i < 64; i++) response.write(Buffer.alloc(32 * 1024)); + await new Promise(resolve => setImmediate(resolve)); + + expect(bytes()).to.equal(2 * 1024 * 1024); + expect(tracking.readableLength).to.equal(0); + expect(tracking.writableLength).to.equal(0); + }); + + it("preserves complete bodies when capture is requested", async () => { + const { response } = createResponse(true); + const body = response.body.asBuffer(); + response.write('first '); + response.end('second'); + expect((await body).toString()).to.equal('first second'); + }); + + it("still drops oversized captured bodies without queuing later chunks", async () => { + const { response, tracking } = createResponse(true, 4); + const body = response.body.asBuffer(); + response.write('too large'); + for (let i = 0; i < 64; i++) response.write(Buffer.alloc(32 * 1024)); + response.end(); + expect((await body).length).to.equal(0); + await new Promise(resolve => setImmediate(resolve)); + expect(tracking.readableLength).to.equal(0); + expect(tracking.writableLength).to.equal(0); + }); + + it("emits body data without retaining a complete response", async () => { + const finished = getDeferred(); + const chunks: Buffer[] = []; + const { response, tracking } = createResponse(false, 4, (_id, _time, content, ended) => { + chunks.push(Buffer.from(content)); + if (ended) finished.resolve(); + }); + response.write('first '); + response.end('second'); + await finished; + expect(Buffer.concat(chunks).toString()).to.equal('first second'); + expect(tracking.readableLength).to.equal(0); + expect(tracking.writableLength).to.equal(0); + }); + }); +});