Skip to content

Commit f4c11cc

Browse files
xia-chaoaduh95
authored andcommitted
zlib: reject reset while a zstd frame is incomplete
Resetting a ZstdCompress stream while a frame is still in progress dropped the frame state, but any bytes already written out stayed at the start of the output stream. The next frame was then appended to that fragment, so the resulting stream could not be decompressed. The failure was silent: the compressor reported no error at all. ZSTD_reset_session_only cancels unflushed internal data, so refuse reset only when the incomplete frame has already emitted output. Signed-off-by: Xia Chao <shapirolutts@gmail.com> PR-URL: #66088 Reviewed-By: James M Snell <jasnell@gmail.com>
1 parent ae0d141 commit f4c11cc

3 files changed

Lines changed: 175 additions & 8 deletions

File tree

‎doc/api/zlib.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2196,6 +2196,12 @@ configured for a Zstd compressor, it applies again to the next frame.
21962196

21972197
Calling `reset()` while a write is in progress throws an `Error`.
21982198

2199+
Resetting an incomplete Zstd compression frame after it has emitted output
2200+
causes the stream to error with `ERR_ZLIB_INCOMPLETE_FRAME`. Resetting at
2201+
that point would discard the frame state while the bytes already written
2202+
out remain at the start of the output stream, leaving it undecodable. Call
2203+
`.end()`, or start over with a new stream, instead.
2204+
21992205
## Class: `ZstdOptions`
22002206

22012207
> Stability: 1 - Experimental

‎src/node_zlib.cc‎

Lines changed: 40 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -352,6 +352,14 @@ class ZstdCompressContext final : public ZstdContext {
352352

353353
uint64_t pledged_src_size_ = ZSTD_CONTENTSIZE_UNKNOWN;
354354
std::optional<uint64_t> consumed_src_size_;
355+
356+
// A frame is complete once ZSTD_compressStream2() has been called with
357+
// ZSTD_e_end and has returned 0. Resetting an incomplete frame is only unsafe
358+
// once some of that frame has already been written out: those bytes cannot be
359+
// discarded, and the next frame would be appended to the fragment. Unflushed
360+
// internal state alone is cancelled by ZSTD_reset_session_only.
361+
bool frame_complete_ = true;
362+
bool frame_output_emitted_ = false;
355363
};
356364

357365
class ZstdDecompressContext final : public ZstdContext {
@@ -1678,6 +1686,8 @@ CompressionError ZstdCompressContext::Init(uint64_t pledged_src_size,
16781686
std::string_view dictionary,
16791687
bool) {
16801688
pledged_src_size_ = pledged_src_size;
1689+
frame_complete_ = true;
1690+
frame_output_emitted_ = false;
16811691
if (pledged_src_size == ZSTD_CONTENTSIZE_UNKNOWN) {
16821692
consumed_src_size_.reset();
16831693
} else {
@@ -1718,6 +1728,17 @@ CompressionError ZstdCompressContext::Init(uint64_t pledged_src_size,
17181728
}
17191729

17201730
CompressionError ZstdCompressContext::ResetStream() {
1731+
// ZSTD_reset_session_only cancels unflushed internal data. Bytes that have
1732+
// already been written out cannot be taken back, so refuse reset only when
1733+
// the current incomplete frame has already emitted output.
1734+
if (!frame_complete_ && frame_output_emitted_) {
1735+
return CompressionError(
1736+
"Cannot reset a zstd stream with an incomplete frame; end the frame "
1737+
"or discard the output produced so far",
1738+
"ERR_ZLIB_INCOMPLETE_FRAME",
1739+
ZSTD_error_stage_wrong);
1740+
}
1741+
17211742
size_t result = ZSTD_CCtx_reset(cctx_.get(), ZSTD_reset_session_only);
17221743
if (ZSTD_isError(result)) {
17231744
const ZSTD_ErrorCode error = ZSTD_getErrorCode(result);
@@ -1737,6 +1758,8 @@ CompressionError ZstdCompressContext::ResetStream() {
17371758
} else {
17381759
consumed_src_size_ = 0;
17391760
}
1761+
frame_complete_ = true;
1762+
frame_output_emitted_ = false;
17401763
error_ = ZSTD_error_no_error;
17411764
error_string_.clear();
17421765
error_code_string_.clear();
@@ -1751,19 +1774,28 @@ void ZstdCompressContext::DoThreadPoolWork() {
17511774
if (consumed_src_size_.has_value()) {
17521775
*consumed_src_size_ += input_.pos - input_pos;
17531776
}
1777+
if (output_.pos > 0) {
1778+
frame_output_emitted_ = true;
1779+
}
17541780
if (ZSTD_isError(remaining)) {
17551781
error_ = ZSTD_getErrorCode(remaining);
17561782
error_code_string_ = ZstdStrerror(error_);
17571783
error_string_ = ZSTD_getErrorString(error_);
1758-
} else if (remaining == 0 && flush_ == ZSTD_e_end &&
1759-
consumed_src_size_.has_value()) {
1760-
uint64_t const consumed_src_size = *consumed_src_size_;
1761-
consumed_src_size_.reset();
1762-
if (consumed_src_size != pledged_src_size_) {
1763-
error_ = ZSTD_error_srcSize_wrong;
1764-
error_code_string_ = ZstdStrerror(error_);
1765-
error_string_ = ZSTD_getErrorString(error_);
1784+
frame_complete_ = false;
1785+
} else if (remaining == 0 && flush_ == ZSTD_e_end) {
1786+
frame_complete_ = true;
1787+
frame_output_emitted_ = false;
1788+
if (consumed_src_size_.has_value()) {
1789+
uint64_t const consumed_src_size = *consumed_src_size_;
1790+
consumed_src_size_.reset();
1791+
if (consumed_src_size != pledged_src_size_) {
1792+
error_ = ZSTD_error_srcSize_wrong;
1793+
error_code_string_ = ZstdStrerror(error_);
1794+
error_string_ = ZSTD_getErrorString(error_);
1795+
}
17661796
}
1797+
} else {
1798+
frame_complete_ = false;
17671799
}
17681800
}
17691801

Lines changed: 129 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,129 @@
1+
'use strict';
2+
3+
// Tests that reset() refuses to run only when an incomplete zstd frame has
4+
// already emitted output.
5+
//
6+
// ZSTD_reset_session_only cancels unflushed internal data, so write()-then-
7+
// reset() with no emitted bytes is safe. Once flush() (or a write that filled
8+
// the output buffer) has written a fragment out, resetting would append the
9+
// next frame to that fragment and produce an undecodable stream.
10+
11+
require('../common');
12+
const assert = require('assert');
13+
const { finished } = require('stream/promises');
14+
const test = require('node:test');
15+
const zlib = require('zlib');
16+
17+
test('ZstdCompress reset throws when an incomplete frame has emitted output',
18+
async () => {
19+
const stream = zlib.createZstdCompress();
20+
const chunks = [];
21+
stream.on('data', (chunk) => chunks.push(chunk));
22+
23+
stream.write(Buffer.from('hello'));
24+
await new Promise((resolve) => stream.flush(resolve));
25+
assert.ok(Buffer.concat(chunks).length > 0);
26+
27+
// A fragment of the frame is already outside the compressor, so reset()
28+
// must refuse instead of silently producing a stream that cannot be
29+
// decoded.
30+
stream.reset();
31+
stream.end(Buffer.from('world'));
32+
33+
await assert.rejects(finished(stream), {
34+
code: 'ERR_ZLIB_INCOMPLETE_FRAME',
35+
});
36+
});
37+
38+
test('ZstdCompress reset throws when write itself emitted frame output',
39+
async () => {
40+
const stream = zlib.createZstdCompress();
41+
const chunks = [];
42+
stream.on('data', (chunk) => chunks.push(chunk));
43+
44+
// Fill a buffer that is large enough for zstd to emit compressed bytes
45+
// during write(), without an explicit flush().
46+
const input = Buffer.allocUnsafe(512 * 1024);
47+
for (let i = 0; i < input.length; i++) {
48+
input[i] = i & 0xff;
49+
}
50+
await new Promise((resolve, reject) => {
51+
stream.write(input, (err) => {
52+
if (err) {
53+
reject(err);
54+
} else {
55+
resolve();
56+
}
57+
});
58+
});
59+
assert.ok(Buffer.concat(chunks).length > 0);
60+
61+
stream.reset();
62+
stream.end(Buffer.from('world'));
63+
64+
await assert.rejects(finished(stream), {
65+
code: 'ERR_ZLIB_INCOMPLETE_FRAME',
66+
});
67+
});
68+
69+
test('ZstdCompress reset after write without emitted output still works',
70+
async () => {
71+
const stream = zlib.createZstdCompress();
72+
const chunks = [];
73+
stream.on('data', (chunk) => chunks.push(chunk));
74+
75+
// Small writes often stay buffered inside zstd until flush/end.
76+
await new Promise((resolve, reject) => {
77+
stream.write(Buffer.from('hello'), (err) => {
78+
if (err) {
79+
reject(err);
80+
} else {
81+
resolve();
82+
}
83+
});
84+
});
85+
assert.strictEqual(Buffer.concat(chunks).length, 0);
86+
87+
// No bytes have left the compressor, so session reset is allowed.
88+
stream.reset();
89+
stream.end(Buffer.from('world'));
90+
await finished(stream);
91+
92+
assert.strictEqual(
93+
zlib.zstdDecompressSync(Buffer.concat(chunks)).toString(),
94+
'world',
95+
);
96+
});
97+
98+
test('ZstdCompress flush followed by end still produces a valid stream',
99+
async () => {
100+
const stream = zlib.createZstdCompress();
101+
const chunks = [];
102+
stream.on('data', (chunk) => chunks.push(chunk));
103+
104+
stream.write(Buffer.from('hello'));
105+
await new Promise((resolve) => stream.flush(resolve));
106+
stream.end(Buffer.from('world'));
107+
await finished(stream);
108+
109+
assert.strictEqual(
110+
zlib.zstdDecompressSync(Buffer.concat(chunks)).toString(),
111+
'helloworld',
112+
);
113+
});
114+
115+
test('ZstdCompress reset before any write still works', async () => {
116+
const stream = zlib.createZstdCompress();
117+
const chunks = [];
118+
stream.on('data', (chunk) => chunks.push(chunk));
119+
120+
// No frame has been started yet, so reset() is allowed.
121+
stream.reset();
122+
stream.end(Buffer.from('hello'));
123+
await finished(stream);
124+
125+
assert.strictEqual(
126+
zlib.zstdDecompressSync(Buffer.concat(chunks)).toString(),
127+
'hello',
128+
);
129+
});

0 commit comments

Comments
 (0)