Skip to content

Commit f59b0f2

Browse files
authored
fix(readable): keep body bytes that arrive after setEncoding() (#5620)
* fix(readable): keep body bytes that arrive after setEncoding() A response body was silently truncated when setEncoding() was used and the body arrived in more than one chunk: .text(), .json(), .bytes(), .arrayBuffer() and .blob() all returned short data and raised nothing. setEncoding() installs a StringDecoder, so state.buffer holds decoded strings from then on and the trailing bytes of a multi-byte sequence split across a chunk boundary sit inside the decoder. #5003 worked around that by snapshotting the raw chunks buffered at the time of the call into kPreservedBuffer, but chunks arriving afterwards were never added, and consumeStart() treated the two sources as mutually exclusive, so everything but the snapshot was dropped. Drop kPreservedBuffer and read the single source of truth instead: state.buffer, re-encoded back to bytes by consumePush(), plus the bytes the decoder is still holding. Nothing is retained beside the stream's own buffer, and push() stays untouched on the hot path. Fixes: #5611 * test: drive the setEncoding() regression test off events, not timers The chunk that has to arrive after setEncoding() is now released by a handshake with the server rather than a 50ms timer, and the wait for the body to arrive is the client's 'drain' event rather than a 150ms sleep.
1 parent f6bfa14 commit f59b0f2

3 files changed

Lines changed: 201 additions & 38 deletions

File tree

‎lib/api/readable.js‎

Lines changed: 23 additions & 38 deletions
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,6 @@ const kContentType = Symbol('kContentType')
1515
const kContentLength = Symbol('kContentLength')
1616
const kUsed = Symbol('kUsed')
1717
const kBytesRead = Symbol('kBytesRead')
18-
const kPreservedBuffer = Symbol('kPreservedBuffer')
1918

2019
const noop = () => {}
2120

@@ -326,36 +325,14 @@ class BodyReadable extends Readable {
326325
*/
327326
setEncoding (encoding) {
328327
if (Buffer.isEncoding(encoding)) {
329-
// Preserve raw Buffer chunks for the consume path (body.text(),
330-
// body.json(), etc.) before super.setEncoding() replaces them
331-
// with decoded strings. Without this, the consume path would
332-
// lose access to the original bytes — some of which may be held
333-
// by the decoder for incomplete multi-byte sequences, and the
334-
// rest converted to strings that can't be safely concatenated
335-
// byte-wise.
336-
const state = this._readableState
337-
const buffer = state.buffer
338-
if (buffer && state.length > 0) {
339-
const bufferIndex = state.bufferIndex ?? 0
340-
const preserved = []
341-
const source = typeof buffer.slice === 'function'
342-
? buffer.slice(bufferIndex)
343-
: buffer
344-
for (const data of source) {
345-
if (Buffer.isBuffer(data)) {
346-
preserved.push(data)
347-
}
348-
}
349-
if (preserved.length > 0) {
350-
this[kPreservedBuffer] = (this[kPreservedBuffer] || []).concat(preserved)
351-
}
352-
}
353-
354328
// Delegate to Node.js Readable.setEncoding() which initializes a
355329
// StringDecoder and re-encodes already-buffered chunks. This properly
356330
// handles multi-byte sequences split at chunk boundaries for the
357331
// for-await / on('data') paths. Without this, Node.js uses
358332
// buf.toString(encoding) on each chunk, producing U+FFFD for split chars.
333+
//
334+
// The consume path (body.text(), body.json(), ...) copes with the
335+
// decoded strings this leaves in state.buffer, see consumeStart().
359336
super.setEncoding(encoding)
360337
}
361338
return this
@@ -464,17 +441,7 @@ function consumeStart (consume) {
464441

465442
const { _readableState: state } = consume.stream
466443

467-
// If setEncoding() was called, state.buffer may contain decoded strings
468-
// (which would break Buffer.concat in chunksDecode). Use the preserved
469-
// raw Buffers (saved before super.setEncoding() in setEncoding()) for
470-
// byte-level accurate consumption. Otherwise read from state.buffer.
471-
const preserved = consume.stream[kPreservedBuffer]
472-
if (preserved && preserved.length > 0) {
473-
for (const chunk of preserved) {
474-
consumePush(consume, chunk)
475-
}
476-
consume.stream[kPreservedBuffer] = null
477-
} else if (state.bufferIndex) {
444+
if (state.bufferIndex) {
478445
const start = state.bufferIndex
479446
const end = state.buffer.length
480447
for (let n = start; n < end; n++) {
@@ -486,6 +453,16 @@ function consumeStart (consume) {
486453
}
487454
}
488455

456+
// If setEncoding() was called, state.buffer holds decoded strings, which
457+
// consumePush() turns back into bytes. The trailing bytes of a multi-byte
458+
// sequence split across a chunk boundary are not part of any of those
459+
// strings, they are held inside the decoder until the rest arrives, so
460+
// take them from there.
461+
const decoder = state.decoder
462+
if (decoder != null && decoder.lastNeed > 0) {
463+
consumePush(consume, Buffer.from(decoder.lastChar.subarray(0, decoder.lastTotal - decoder.lastNeed)))
464+
}
465+
489466
if (state.endEmitted) {
490467
// No `this` to read the consume off here: consumeStart is a free function, called from
491468
// the queueMicrotask above. The callback below does have one, because the emitter passes
@@ -588,14 +565,22 @@ function consumeEnd (consume, encoding) {
588565

589566
/**
590567
* @param {Consume} consume
591-
* @param {Buffer} chunk
568+
* @param {Buffer|string} chunk
592569
* @returns {void}
593570
*/
594571
function consumePush (consume, chunk) {
595572
if (consume.body === null) {
596573
return
597574
}
598575

576+
if (typeof chunk === 'string') {
577+
// Buffered before the consume started, while an encoding was set.
578+
// consume.length has to stay a byte count and chunksDecode()/chunksConcat()
579+
// only work on bytes, so re-encode. A string's own length is in UTF-16 code
580+
// units and Uint8Array.prototype.set() ignores a string argument entirely.
581+
chunk = Buffer.from(chunk, consume.stream._readableState.encoding)
582+
}
583+
599584
consume.length += chunk.length
600585
consume.body.push(chunk)
601586
}

‎test/client-request.js‎

Lines changed: 46 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1574,6 +1574,52 @@ test('setEncoding(\'utf8\') handles 3-byte UTF-8 characters split across chunks'
15741574
await t.completed
15751575
})
15761576

1577+
test('#5611 - setEncoding() then .text() does not truncate a body that arrives in several chunks', async (t) => {
1578+
t = tspl(t, { plan: 2 })
1579+
1580+
// '傳' is 3 bytes in UTF-8, split across the two chunks below.
1581+
const text = 'abc傳def'
1582+
const buf = Buffer.from(text)
1583+
1584+
// The rest of the body is released only once the client has called
1585+
// setEncoding(), which is what puts the second chunk on the wrong side of it.
1586+
const sendRest = new EE()
1587+
1588+
const server = createServer({ joinDuplicateHeaders: true }, async (req, res) => {
1589+
res.writeHead(200, { 'content-type': 'text/plain; charset=utf-8' })
1590+
res.write(buf.subarray(0, 4))
1591+
await EE.once(sendRest, 'go')
1592+
res.end(buf.subarray(4))
1593+
})
1594+
after(() => {
1595+
server.closeAllConnections?.()
1596+
server.close()
1597+
})
1598+
1599+
server.listen(0, async () => {
1600+
const client = new Client(`http://localhost:${server.address().port}`)
1601+
after(client.destroy.bind(client))
1602+
1603+
const { body } = await client.request({ path: '/', method: 'GET' })
1604+
body.setEncoding('utf8')
1605+
1606+
// 'drain' means the request has run to completion, so the whole body has
1607+
// arrived and is buffered as decoded strings, with the split character
1608+
// held inside the decoder. Waiting on the body itself is not an option:
1609+
// reading it would disturb it and make the consume below unusable.
1610+
const completed = EE.once(client, 'drain')
1611+
sendRest.emit('go')
1612+
await completed
1613+
1614+
const result = await body.text()
1615+
1616+
t.strictEqual(result, text)
1617+
t.strictEqual(Buffer.byteLength(result), buf.length)
1618+
})
1619+
1620+
await t.completed
1621+
})
1622+
15771623
test('#3736 - Aborted Response (without consuming body)', async (t) => {
15781624
const plan = tspl(t, { plan: 1 })
15791625

‎test/readable.js‎

Lines changed: 132 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -226,4 +226,136 @@ describe('Readable', () => {
226226

227227
t.strictEqual(text, 'hello world')
228228
})
229+
230+
test('setEncoding() then .text() keeps chunks pushed after setEncoding()', async function (t) {
231+
t = tspl(t, { plan: 2 })
232+
233+
function resume () {
234+
}
235+
function abort () {
236+
}
237+
const r = new Readable({ resume, abort })
238+
239+
// '傳' is 3 bytes in UTF-8, so cutting every 2 bytes splits each of them
240+
// across a chunk boundary.
241+
const expected = 'a傳b傳c傳d'
242+
const buf = Buffer.from(expected)
243+
const chunks = []
244+
for (let n = 0; n < buf.length; n += 2) {
245+
chunks.push(buf.subarray(n, n + 2))
246+
}
247+
248+
// Buffered when setEncoding() runs: these are replaced by a single decoded
249+
// string, with the tail of the split '傳' held inside the decoder.
250+
r.push(chunks[0])
251+
r.push(chunks[1])
252+
253+
r.setEncoding('utf8')
254+
255+
// Pushed after setEncoding() but before .text() is called: these are
256+
// buffered as decoded strings, not as bytes.
257+
r.push(chunks[2])
258+
r.push(chunks[3])
259+
260+
const promise = r.text()
261+
262+
// Pushed after .text() but before the consume actually starts.
263+
r.push(chunks[4])
264+
265+
setImmediate(() => {
266+
// Pushed once the consume is running.
267+
r.push(chunks[5])
268+
r.push(chunks[6])
269+
r.push(null)
270+
})
271+
272+
const text = await promise
273+
274+
t.strictEqual(text, expected)
275+
t.strictEqual(Buffer.byteLength(text), buf.length)
276+
})
277+
278+
test('setEncoding() with only a partial character buffered', async function (t) {
279+
t = tspl(t, { plan: 1 })
280+
281+
function resume () {
282+
}
283+
function abort () {
284+
}
285+
const r = new Readable({ resume, abort })
286+
287+
const buf = Buffer.from('傳')
288+
289+
// Decodes to the empty string, so setEncoding() leaves nothing buffered
290+
// and the byte stays inside the decoder.
291+
r.push(buf.subarray(0, 1))
292+
293+
r.setEncoding('utf8')
294+
295+
r.push(buf.subarray(1))
296+
297+
process.nextTick(() => {
298+
r.push(null)
299+
})
300+
301+
t.strictEqual(await r.text(), '傳')
302+
})
303+
304+
for (const encoding of ['utf8', 'hex', 'base64', 'latin1']) {
305+
test(`setEncoding('${encoding}') before any chunk arrives`, async function (t) {
306+
t = tspl(t, { plan: 5 })
307+
308+
function resume () {
309+
}
310+
function abort () {
311+
}
312+
313+
const buf = Buffer.from('hello 傳 world')
314+
315+
// Nothing is buffered when setEncoding() runs, so every chunk reaches
316+
// state.buffer as a decoded string.
317+
function body () {
318+
const r = new Readable({ resume, abort })
319+
r.setEncoding(encoding)
320+
r.push(buf.subarray(0, 4))
321+
r.push(buf.subarray(4, 8))
322+
process.nextTick(() => {
323+
r.push(buf.subarray(8))
324+
r.push(null)
325+
})
326+
return r
327+
}
328+
329+
t.deepStrictEqual(await body().bytes(), new Uint8Array(buf))
330+
t.deepStrictEqual(new Uint8Array(await body().arrayBuffer()), new Uint8Array(buf))
331+
332+
const blob = await body().blob()
333+
t.strictEqual(blob.size, buf.length)
334+
t.deepStrictEqual(Buffer.from(await blob.arrayBuffer()), buf)
335+
336+
t.strictEqual(await body().text(), buf.toString(encoding))
337+
})
338+
}
339+
340+
test('setEncoding() then .json()', async function (t) {
341+
t = tspl(t, { plan: 1 })
342+
343+
function resume () {
344+
}
345+
function abort () {
346+
}
347+
const r = new Readable({ resume, abort })
348+
349+
const buf = Buffer.from(JSON.stringify({ hello: '傳' }))
350+
351+
r.setEncoding('utf8')
352+
r.push(buf.subarray(0, 12))
353+
354+
process.nextTick(() => {
355+
r.push(buf.subarray(12))
356+
r.push(null)
357+
})
358+
359+
t.deepStrictEqual(await r.json(), { hello: '傳' })
360+
})
229361
})

0 commit comments

Comments
 (0)