Skip to content

Commit 2746337

Browse files
committed
Extract body buffering into stream-capture (with small fixes en route)
1 parent 7acf136 commit 2746337

7 files changed

Lines changed: 569 additions & 150 deletions

File tree

‎src/rules/requests/request-step-impls.ts‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1524,7 +1524,9 @@ export class WaitForRequestBodyStepImpl extends WaitForRequestBodyStep {
15241524
static readonly fromDefinition = () => new WaitForRequestBodyStepImpl();
15251525

15261526
async handle(request: OngoingRequest): Promise<{ continue: true }> {
1527-
await request.body.asBuffer();
1527+
// Wait for body delivery to complete - not just buffering to end (they can be
1528+
// different, e.g. if buffer tops out the max buffering limits).
1529+
await request.body.waitForEnd();
15281530
return { continue: true };
15291531
}
15301532
}

‎src/types.ts‎

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -214,7 +214,15 @@ export interface OngoingRequest extends Request, stream.Readable {
214214

215215
export interface OngoingBody {
216216
asStream: () => stream.Readable;
217+
/**
218+
* Resolves with the body content when it's done, or an empty body if the body is dropped
219+
* due to hitting the max buffering limits.
220+
*/
217221
asBuffer: () => Promise<Buffer>;
222+
/**
223+
* Resolves at the end of the body stream - regardless of buffering behaviour.
224+
*/
225+
waitForEnd: () => Promise<void>;
218226
asDecodedBuffer: () => Promise<Buffer>;
219227
asText: () => Promise<string>;
220228
asJson: () => Promise<object>;

‎src/util/buffer-utils.ts‎

Lines changed: 17 additions & 127 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,6 @@
11
import { Buffer } from 'buffer';
2-
import { EventEmitter } from 'events';
32
import * as stream from 'stream';
43

5-
import { isNode } from './util';
6-
7-
const MAX_BUFFER_SIZE = isNode
8-
? require('buffer').constants.MAX_LENGTH
9-
: Infinity;
10-
114
export const asBuffer = (input: Buffer | Uint8Array | string) =>
125
Buffer.isBuffer(input)
136
? input
@@ -16,132 +9,29 @@ export const asBuffer = (input: Buffer | Uint8Array | string) =>
169
// Is Array:
1710
: Buffer.from(input);
1811

19-
export type BufferInProgress = Promise<Buffer> & {
20-
currentChunks: Buffer[]; // Stores the body chunks as they arrive
21-
failedWith?: Error; // Stores the error that killed the stream, if one did
22-
events: EventEmitter; // Emits events - notably 'truncate' if data is truncated
23-
};
24-
25-
// Takes a buffer and a stream, returns a simple stream that outputs the buffer then the stream. The stream
26-
// is lazy, so doesn't read data in from the buffer or input until something here starts reading.
27-
export const bufferThenStream = (buffer: BufferInProgress, inputStream: stream.Readable): stream.Readable => {
28-
let active = false;
29-
30-
const outputStream = new stream.PassThrough({
31-
// Note we use the default highWaterMark, which means this applies backpressure, pushing buffering
32-
// onto the OS & backpressure on network instead of accepting data before we're ready to stream it.
33-
34-
// Without changing behaviour, we listen for read() events, and don't start streaming until we get one.
35-
read(size) {
36-
// On the first actual read of this stream, we pull from the buffer
37-
// and then hook it up to the input.
38-
if (!active) {
39-
if (buffer.failedWith) {
40-
outputStream.destroy(buffer.failedWith);
41-
} else {
42-
// First stream everything that's been buffered so far:
43-
outputStream.write(Buffer.concat(buffer.currentChunks));
44-
45-
// Then start streaming all future incoming data:
46-
inputStream.pipe(outputStream);
47-
48-
if (inputStream.readableEnded) outputStream.end();
49-
if (inputStream.readableAborted) outputStream.destroy();
50-
51-
// Forward any future errors from the input stream:
52-
inputStream.on('error', (e) => {
53-
outputStream.emit('error', e)
54-
});
55-
56-
// Silence 'unhandled rejection' warnings here, since we'll handle
57-
// them on the stream instead
58-
buffer.catch(() => {});
59-
}
60-
active = true;
61-
}
62-
63-
// Except for the first activation logic (above) do the default transform read() steps just
64-
// like a normal PassThrough stream.
65-
return stream.Transform.prototype._read.call(this, size);
66-
}
67-
});
68-
69-
buffer.events.on('truncate', (chunks) => {
70-
// If the stream hasn't started yet, start it now, so it grabs the buffer
71-
// data before it gets truncated:
72-
if (!active) outputStream.read(0);
73-
});
74-
75-
return outputStream;
76-
};
77-
7812
export const bufferToStream = (buffer: Buffer): stream.Readable => {
7913
const outputStream = new stream.PassThrough();
8014
outputStream.end(buffer);
8115
return outputStream;
8216
};
8317

84-
export const streamToBuffer = (input: stream.Readable, maxSize = MAX_BUFFER_SIZE) => {
85-
let chunks: Buffer[] = [];
86-
87-
const bufferPromise = <BufferInProgress> new Promise(
88-
(resolve, reject) => {
89-
function failWithAbortError() {
90-
bufferPromise.failedWith = new Error('Aborted');
91-
reject(bufferPromise.failedWith);
92-
}
93-
94-
// If stream has already finished/aborted, resolve accordingly immediately:
95-
if (input.readableEnded) return resolve(Buffer.from([]));
96-
if (input.readableAborted) return setImmediate(failWithAbortError);
97-
98-
let currentSize = 0;
99-
const onData = (d: Buffer) => {
100-
currentSize += d.length;
101-
chunks.push(d);
102-
103-
// If we go over maxSize, drop the whole stream, so the buffer
104-
// resolves empty. MaxSize should be large, so this is rare,
105-
// and only happens as an alternative to crashing the process.
106-
if (currentSize > maxSize) {
107-
// Announce truncation, so that other mechanisms (asStream) can
108-
// capture this data if they're interested in it.
109-
bufferPromise.events.emit('truncate', chunks);
110-
111-
// Drop all the data so far & stop reading
112-
bufferPromise.currentChunks = chunks = [];
113-
input.removeListener('data', onData);
114-
115-
// We then resolve immediately - the buffer is done, even if the body
116-
// might still be streaming in we're not listening to it. This means
117-
// that requests can 'complete' for event/callback purposes while
118-
// they're actually still streaming, but only in this scenario where
119-
// the data is too large to really be used by the events/callbacks.
120-
121-
// If we don't resolve, then cases which intentionally don't consume
122-
// the raw stream but do consume the buffer (beforeRequest) would
123-
// deadlock: beforeRequest must complete to begin streaming the
124-
// full body to the target clients.
125-
126-
resolve(Buffer.from([]));
127-
}
128-
};
129-
input.on('data', onData);
130-
131-
input.once('end', () => {
132-
resolve(Buffer.concat(chunks));
133-
});
134-
input.once('aborted', failWithAbortError);
135-
input.on('error', (e) => {
136-
bufferPromise.failedWith = bufferPromise.failedWith || e;
137-
reject(e);
138-
});
139-
}
140-
);
141-
bufferPromise.currentChunks = chunks;
142-
bufferPromise.events = new EventEmitter();
143-
return bufferPromise;
144-
};
18+
/**
19+
* Reads a stream to completion, returning all its data as a single buffer. Rejects if the
20+
* stream errors or is aborted.
21+
*/
22+
export const streamToBuffer = (input: stream.Readable) => new Promise<Buffer>((resolve, reject) => {
23+
const failWithAbortError = () => reject(new Error('Aborted'));
24+
25+
// If the stream has already finished/aborted, resolve accordingly immediately:
26+
if (input.readableEnded) return resolve(Buffer.from([]));
27+
if (input.readableAborted) return setImmediate(failWithAbortError);
28+
29+
const chunks: Buffer[] = [];
30+
input.on('data', (d: Buffer) => chunks.push(d));
31+
input.once('end', () => resolve(Buffer.concat(chunks)));
32+
input.once('aborted', failWithAbortError);
33+
input.on('error', reject);
34+
});
14535

14636
export function splitBuffer(input: Buffer, splitter: string, maxParts = Infinity) {
14737
const parts: Buffer[] = [];

‎src/util/request-utils.ts‎

Lines changed: 14 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -33,13 +33,10 @@ import {
3333
import { makePropertyWritable } from './util';
3434
import { Mutable } from './type-utils';
3535
import {
36-
bufferThenStream,
37-
bufferToStream,
38-
BufferInProgress,
3936
splitBuffer,
40-
streamToBuffer,
4137
asBuffer
4238
} from './buffer-utils';
39+
import { StreamCapture } from './stream-capture';
4340
import {
4441
flattenPairedRawHeaders,
4542
getHeaderValue,
@@ -145,33 +142,28 @@ export async function decodeBodyBuffer(buffer: Buffer, headers: Headers) {
145142
// Parse an in-progress request or response stream, i.e. where the body or possibly even the headers have
146143
// not been fully received/sent yet.
147144
const parseBodyStream = (bodyStream: stream.Readable, maxSize: number, getHeaders: () => Headers): OngoingBody => {
148-
let bufferPromise: BufferInProgress | null = null;
149-
let completedBuffer: Buffer | null = null;
145+
let capture: StreamCapture | null = null;
146+
147+
const startCapture = () => {
148+
if (!capture) {
149+
capture = new StreamCapture(bodyStream, maxSize);
150+
}
151+
return capture;
152+
};
150153

151154
let body = {
152155
// Returns a stream for the full body, not the live streaming body.
153156
// Each call creates a new stream, which starts with the already seen
154157
// and buffered data, and then continues with the live stream, if active.
155158
// Listeners to this stream *must* be attached synchronously after this call.
156159
asStream() {
157-
// If we've already buffered the whole body, just stream it out:
158-
if (completedBuffer) return bufferToStream(completedBuffer);
159-
160-
// Otherwise, we want to start buffering now, and wrap that with
161-
// a stream that can live-stream the buffered data on demand:
162-
const buffer = body.asBuffer();
163-
buffer.catch(() => {}); // Errors will be handled via the stream, so silence unhandled rejections here.
164-
return bufferThenStream(buffer, bodyStream);
160+
return startCapture().takeStream();
165161
},
166162
asBuffer() {
167-
if (!bufferPromise) {
168-
bufferPromise = streamToBuffer(bodyStream, maxSize);
169-
170-
bufferPromise
171-
.then((buffer) => completedBuffer = buffer)
172-
.catch(() => {}); // If we get no body, completedBuffer stays null
173-
}
174-
return bufferPromise;
163+
return startCapture().buffer;
164+
},
165+
waitForEnd() {
166+
return startCapture().completed;
175167
},
176168
async asDecodedBuffer() {
177169
const buffer = await body.asBuffer();

0 commit comments

Comments
 (0)