Skip to content
File

Blob: src/workerd/api/node/zlib-util.c++

44.7 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4// Copyright Joyent and Node contributors. All rights reserved. MIT license.
5 
6#include "zlib-util.h"
7 
8#include "util.h"
9 
10// The following implementation is adapted from Node.js
11// and therefore follows Node.js style as opposed to kj style.
12// Latest implementation of Node.js zlib can be found at:
13// https://github.com/nodejs/node/blob/main/src/node_zlib.cc
14namespace workerd::api::node {
15 
16kj::ArrayPtr<const kj::byte> getInputFromSource(const ZlibUtil::InputSource& data) {
17 KJ_SWITCH_ONEOF(data) {
18 KJ_CASE_ONEOF(dataBuf, kj::Array<kj::byte>) {
19 JSG_REQUIRE(dataBuf.size() < Z_MAX_CHUNK, RangeError, "Memory limit exceeded"_kj);
20 return dataBuf.asPtr();
21 }
22 
23 KJ_CASE_ONEOF(dataStr, jsg::NonCoercible<kj::String>) {
24 JSG_REQUIRE(dataStr.value.size() < Z_MAX_CHUNK, RangeError, "Memory limit exceeded"_kj);
25 return dataStr.value.asBytes();
26 }
27 }
28 
29 KJ_UNREACHABLE;
30}
31 
32uint32_t ZlibUtil::crc32Sync(InputSource data, uint32_t value) {
33 auto dataPtr = getInputFromSource(data);
34 return crc32(value, dataPtr.begin(), dataPtr.size());
35}
36 
37namespace {
38class GrowableBuffer final {
39 // A copy of kj::Vector with some additional methods for use as a growable buffer with a maximum
40 // size
41 public:
42 inline explicit GrowableBuffer(size_t _chunkSize, size_t _maxCapacity)
43 : maxCapacity(_maxCapacity) {
44 auto maxChunkSize = kj::min(_chunkSize, maxCapacity);
45 builder = kj::heapArrayBuilder<kj::byte>(maxChunkSize);
46 chunkSize = maxChunkSize;
47 }
48 
49 size_t size() const {
50 return builder.size();
51 }
52 bool empty() const {
53 return size() == 0;
54 }
55 size_t capacity() const {
56 return builder.capacity();
57 }
58 size_t available() const {
59 return capacity() - size();
60 }
61 
62 kj::byte* begin() KJ_LIFETIMEBOUND {
63 return builder.begin();
64 }
65 kj::byte* end() KJ_LIFETIMEBOUND {
66 return builder.end();
67 }
68 
69 kj::Array<kj::byte> releaseAsArray() {
70 // TODO(perf): Avoid a copy/move by allowing Array<T> to point to incomplete space?
71 if (!builder.isFull()) {
72 setCapacity(size());
73 }
74 return builder.finish();
75 }
76 
77 void adjustUnused(size_t unused) {
78 resize(capacity() - unused);
79 }
80 
81 void resize(size_t size) {
82 if (size > builder.capacity()) grow(size);
83 builder.resize(size);
84 }
85 
86 void addChunk() {
87 reserve(size() + chunkSize);
88 }
89 
90 void reserve(size_t size) {
91 if (size > builder.capacity()) {
92 grow(size);
93 }
94 }
95 
96 bool atMaxCapacity() const {
97 return size() >= maxCapacity;
98 }
99 
100 private:
101 kj::ArrayBuilder<kj::byte> builder;
102 size_t chunkSize;
103 size_t maxCapacity;
104 
105 void grow(size_t minCapacity = 0) {
106 JSG_REQUIRE(minCapacity <= maxCapacity, RangeError, "Memory limit exceeded");
107 setCapacity(kj::min(maxCapacity, kj::max(minCapacity, capacity() == 0 ? 4 : capacity() * 2)));
108 }
109 void setCapacity(size_t newSize) {
110 if (builder.size() > newSize) {
111 builder.truncate(newSize);
112 }
113 
114 kj::ArrayBuilder<kj::byte> newBuilder = kj::heapArrayBuilder<kj::byte>(newSize);
115 newBuilder.addAll(kj::mv(builder));
116 builder = kj::mv(newBuilder);
117 }
118};
119} // namespace
120 
121void ZlibContext::initialize(int _level,
122 int _windowBits,
123 int _memLevel,
124 int _strategy,
125 jsg::Optional<kj::Array<kj::byte>> _dictionary) {
126 if (!((_windowBits == 0) &&
127 (mode == ZlibMode::INFLATE || mode == ZlibMode::GUNZIP || mode == ZlibMode::UNZIP))) {
128 JSG_ASSERT(_windowBits >= Z_MIN_WINDOWBITS && _windowBits <= Z_MAX_WINDOWBITS, RangeError,
129 kj::str("The value of \"options.windowBits\" is out of range. It must be >= ",
130 Z_MIN_WINDOWBITS, " and <= ", Z_MAX_WINDOWBITS, ". Received ", _windowBits));
131 }
132 
133 JSG_REQUIRE(_level >= Z_MIN_LEVEL && _level <= Z_MAX_LEVEL, RangeError,
134 kj::str("The value of \"options.level\" is out of range. It must be >= ", Z_MIN_LEVEL,
135 " and <= ", Z_MAX_LEVEL, ". Received ", _level));
136 JSG_REQUIRE(_memLevel >= Z_MIN_MEMLEVEL && _memLevel <= Z_MAX_MEMLEVEL, RangeError,
137 kj::str("The value of \"options.memLevel\" is out of range. It must be >= ", Z_MIN_MEMLEVEL,
138 " and <= ", Z_MAX_MEMLEVEL, ". Received ", _memLevel));
139 JSG_REQUIRE(_strategy == Z_FILTERED || _strategy == Z_HUFFMAN_ONLY || _strategy == Z_RLE ||
140 _strategy == Z_FIXED || _strategy == Z_DEFAULT_STRATEGY,
141 Error, "invalid strategy"_kj);
142 
143 level = _level;
144 windowBits = _windowBits;
145 memLevel = _memLevel;
146 strategy = _strategy;
147 flush = Z_NO_FLUSH;
148 err = Z_OK;
149 
150 switch (mode) {
151 case ZlibMode::GZIP:
152 case ZlibMode::GUNZIP:
153 windowBits += 16;
154 break;
155 case ZlibMode::UNZIP:
156 windowBits += 32;
157 break;
158 case ZlibMode::DEFLATERAW:
159 case ZlibMode::INFLATERAW:
160 windowBits *= -1;
161 break;
162 default:
163 break;
164 }
165 
166 KJ_IF_SOME(dict, _dictionary) {
167 dictionary = kj::mv(dict);
168 }
169}
170 
171kj::Maybe<CompressionError> ZlibContext::getError() const {
172 // Acceptable error states depend on the type of zlib stream.
173 switch (err) {
174 case Z_OK:
175 case Z_BUF_ERROR:
176 if (stream.avail_out != 0 && flush == Z_FINISH) {
177 return constructError("unexpected end of file"_kj);
178 }
179 break;
180 case Z_STREAM_END:
181 // normal statuses, not fatal
182 break;
183 case Z_NEED_DICT:
184 if (dictionary.empty()) {
185 return constructError("Missing dictionary"_kj);
186 } else {
187 return constructError("Bad dictionary"_kj);
188 }
189 default:
190 // something else.
191 return constructError("Zlib error");
192 }
193 
194 return {};
195}
196 
197kj::Maybe<CompressionError> ZlibContext::setDictionary() {
198 if (dictionary.empty()) {
199 return kj::none;
200 }
201 
202 err = Z_OK;
203 
204 switch (mode) {
205 case ZlibMode::DEFLATE:
206 case ZlibMode::DEFLATERAW:
207 err = deflateSetDictionary(&stream, dictionary.begin(), dictionary.size());
208 break;
209 case ZlibMode::INFLATERAW:
210 err = inflateSetDictionary(&stream, dictionary.begin(), dictionary.size());
211 break;
212 default:
213 break;
214 }
215 
216 if (err != Z_OK) {
217 return constructError("Failed to set dictionary"_kj);
218 }
219 
220 return kj::none;
221}
222 
223bool ZlibContext::initializeZlib() {
224 if (initialized) {
225 return false;
226 }
227 switch (mode) {
228 case ZlibMode::DEFLATE:
229 case ZlibMode::GZIP:
230 case ZlibMode::DEFLATERAW:
231 err = deflateInit2(&stream, level, Z_DEFLATED, windowBits, memLevel, strategy);
232 break;
233 case ZlibMode::INFLATE:
234 case ZlibMode::GUNZIP:
235 case ZlibMode::INFLATERAW:
236 case ZlibMode::UNZIP:
237 err = inflateInit2(&stream, windowBits);
238 break;
239 default:
240 KJ_UNREACHABLE;
241 }
242 
243 if (err != Z_OK) {
244 dictionary.clear();
245 mode = ZlibMode::NONE;
246 return true;
247 }
248 
249 setDictionary();
250 initialized = true;
251 return true;
252}
253 
254kj::Maybe<CompressionError> ZlibContext::resetStream() {
255 bool initialized_now = initializeZlib();
256 if (initialized_now && err != Z_OK) {
257 return constructError("Failed to init stream before reset");
258 }
259 err = Z_OK;
260 switch (mode) {
261 case ZlibMode::DEFLATE:
262 case ZlibMode::DEFLATERAW:
263 case ZlibMode::GZIP:
264 err = deflateReset(&stream);
265 break;
266 case ZlibMode::INFLATE:
267 case ZlibMode::INFLATERAW:
268 case ZlibMode::GUNZIP:
269 err = inflateReset(&stream);
270 break;
271 default:
272 break;
273 }
274 
275 if (err != Z_OK) {
276 return constructError("Failed to reset stream"_kj);
277 }
278 
279 return setDictionary();
280}
281 
282void ZlibContext::work() {
283 bool initialized_now = initializeZlib();
284 if (initialized_now && err != Z_OK) {
285 return;
286 }
287 
288 const Bytef* next_expected_header_byte = nullptr;
289 
290 // If the avail_out is left at 0, then it means that it ran out
291 // of room. If there was avail_out left over, then it means
292 // that all the input was consumed.
293 switch (mode) {
294 case ZlibMode::DEFLATE:
295 case ZlibMode::GZIP:
296 case ZlibMode::DEFLATERAW:
297 err = deflate(&stream, flush);
298 break;
299 case ZlibMode::UNZIP:
300 if (stream.avail_in > 0) {
301 next_expected_header_byte = stream.next_in;
302 }
303 
304 switch (gzip_id_bytes_read) {
305 case 0:
306 if (next_expected_header_byte == nullptr) {
307 break;
308 }
309 
310 if (*next_expected_header_byte == GZIP_HEADER_ID1) {
311 gzip_id_bytes_read = 1;
312 next_expected_header_byte++;
313 
314 if (stream.avail_in == 1) {
315 // The only available byte was already read.
316 break;
317 }
318 } else {
319 mode = ZlibMode::INFLATE;
320 break;
321 }
322 
323 [[fallthrough]];
324 case 1:
325 if (next_expected_header_byte == nullptr) {
326 break;
327 }
328 
329 if (*next_expected_header_byte == GZIP_HEADER_ID2) {
330 gzip_id_bytes_read = 2;
331 mode = ZlibMode::GUNZIP;
332 } else {
333 // There is no actual difference between INFLATE and INFLATERAW
334 // (after initialization).
335 mode = ZlibMode::INFLATE;
336 }
337 
338 break;
339 default:
340 JSG_FAIL_REQUIRE(Error, "Invalid number of gzip magic number bytes read");
341 }
342 
343 [[fallthrough]];
344 case ZlibMode::INFLATE:
345 case ZlibMode::GUNZIP:
346 case ZlibMode::INFLATERAW:
347 err = inflate(&stream, flush);
348 
349 // If data was encoded with dictionary (INFLATERAW will have it set in
350 // SetDictionary, don't repeat that here)
351 if (mode != ZlibMode::INFLATERAW && err == Z_NEED_DICT && !dictionary.empty()) {
352 // Load it
353 err = inflateSetDictionary(&stream, dictionary.begin(), dictionary.size());
354 if (err == Z_OK) {
355 // And try to decode again
356 err = inflate(&stream, flush);
357 } else if (err == Z_DATA_ERROR) {
358 // Both inflateSetDictionary() and inflate() return Z_DATA_ERROR.
359 // Make it possible for After() to tell a bad dictionary from bad
360 // input.
361 err = Z_NEED_DICT;
362 }
363 }
364 
365 while (stream.avail_in > 0 && mode == ZlibMode::GUNZIP && err == Z_STREAM_END &&
366 stream.next_in[0] != 0x00) {
367 // Bytes remain in input buffer. Perhaps this is another compressed
368 // member in the same archive, or just trailing garbage.
369 // Trailing zero bytes are okay, though, since they are frequently
370 // used for padding.
371 
372 resetStream();
373 err = inflate(&stream, flush);
374 }
375 break;
376 default:
377 KJ_UNREACHABLE;
378 }
379}
380 
381kj::Maybe<CompressionError> ZlibContext::setParams(int _level, int _strategy) {
382 bool initialized_now = initializeZlib();
383 if (initialized_now && err != Z_OK) {
384 return constructError("Failed to init stream before set parameters");
385 }
386 err = Z_OK;
387 
388 switch (mode) {
389 case ZlibMode::DEFLATE:
390 case ZlibMode::DEFLATERAW:
391 err = deflateParams(&stream, _level, _strategy);
392 break;
393 default:
394 break;
395 }
396 
397 if (err != Z_OK && err != Z_BUF_ERROR) {
398 return constructError("Failed to set parameters");
399 }
400 
401 return kj::none;
402}
403 
404ZlibContext::~ZlibContext() noexcept(false) {
405 if (!initialized) {
406 return;
407 }
408 
409 auto status = Z_OK;
410 switch (mode) {
411 case ZlibMode::DEFLATE:
412 case ZlibMode::DEFLATERAW:
413 case ZlibMode::GZIP:
414 status = deflateEnd(&stream);
415 break;
416 case ZlibMode::INFLATE:
417 case ZlibMode::INFLATERAW:
418 case ZlibMode::GUNZIP:
419 case ZlibMode::UNZIP:
420 status = inflateEnd(&stream);
421 break;
422 default:
423 break;
424 }
425 
426 JSG_REQUIRE(
427 status == Z_OK || status == Z_DATA_ERROR, Error, "Uncaught error on closing zlib stream");
428}
429 
430void ZlibContext::setBuffers(kj::ArrayPtr<kj::byte> input, kj::ArrayPtr<kj::byte> output) {
431 stream.avail_in = input.size();
432 stream.next_in = input.begin();
433 stream.avail_out = output.size();
434 stream.next_out = output.begin();
435}
436 
437void ZlibContext::setInputBuffer(kj::ArrayPtr<const kj::byte> input) {
438 // The define Z_CONST is not set, so zlib always takes mutable pointers
439 stream.next_in = const_cast<kj::byte*>(input.begin());
440 stream.avail_in = input.size();
441}
442 
443void ZlibContext::setOutputBuffer(kj::ArrayPtr<kj::byte> output) {
444 stream.next_out = output.begin();
445 stream.avail_out = output.size();
446}
447 
448template <typename CompressionContext>
449jsg::Ref<ZlibUtil::CompressionStream<CompressionContext>> ZlibUtil::CompressionStream<
450 CompressionContext>::constructor(jsg::Lock& js, ZlibModeValue mode) {
451 return js.alloc<CompressionStream>(static_cast<ZlibMode>(mode), js.getExternalMemoryTarget());
452}
453 
454template <typename CompressionContext>
455ZlibUtil::CompressionStream<CompressionContext>::~CompressionStream() {
456 JSG_ASSERT(!writing, Error, "Writing to compression stream"_kj);
457 close();
458}
459 
460template <typename CompressionContext>
461void ZlibUtil::CompressionStream<CompressionContext>::emitError(
462 jsg::Lock& js, const CompressionError& error) {
463 KJ_IF_SOME(onError, errorHandler) {
464 onError(js, error.err, kj::mv(error.code), kj::mv(error.message));
465 }
466 
467 writing = false;
468 if (pending_close) {
469 close();
470 }
471}
472 
473template <typename CompressionContext>
474template <bool async>
475void ZlibUtil::CompressionStream<CompressionContext>::writeStream(
476 jsg::Lock& js, int flush, kj::ArrayPtr<kj::byte> input, kj::ArrayPtr<kj::byte> output) {
477 JSG_REQUIRE(initialized, Error, "Writing before initializing"_kj);
478 JSG_REQUIRE(!closed, Error, "Already finalized"_kj);
479 JSG_REQUIRE(!writing, Error, "Writing is in progress"_kj);
480 JSG_REQUIRE(!pending_close, Error, "Pending close"_kj);
481 
482 writing = true;
483 
484 context()->setBuffers(input, output);
485 context()->setFlush(flush);
486 
487 // Clear buffer pointers from the compression context when this scope exits.
488 // The input/output kj::Array parameters are backed by V8 BackingStores
489 // whose lifetimes are tied to their JavaScript ArrayBuffer objects. Once
490 // this method returns, the kj::Array destructor releases its shared_ptr
491 // to the BackingStore. If the JS buffer subsequently becomes unreachable
492 // and is garbage-collected, the BackingStore is freed — leaving any
493 // retained pointers (e.g. z_stream.next_out) dangling. A later call to
494 // deflateParams() (via params()) could then write to freed memory.
495 //
496 // Using KJ_DEFER ensures the pointers are cleared even if an exception
497 // is thrown (e.g. from updateWriteResult() when the writeState buffer is
498 // too small). Without this, the exception would unwind the stack before
499 // reaching an explicit clearBuffers() call, leaving stale pointers.
500 KJ_DEFER(context()->clearBuffers());
501 
502 if constexpr (!async) {
503 context()->work();
504 if (checkError(js)) {
505 writing = false;
506 updateWriteResult(js);
507 }
508 return;
509 }
510 
511 // On Node.js, this is called as a result of `ScheduleWork()` call.
512 // Since, we implement the whole thing as sync, we're going to ahead and call the whole thing here.
513 context()->work();
514 
515 // This is implemented slightly differently in Node.js
516 // Node.js calls AfterThreadPoolWork().
517 // Ref: https://github.com/nodejs/node/blob/9edf4a0856681a7665bd9dcf2ca7cac252784b98/src/node_zlib.cc#L402
518 writing = false;
519 if (!checkError(js)) {
520 return;
521 }
522 updateWriteResult(js);
523 KJ_IF_SOME(cb, writeCallback) {
524 cb(js);
525 }
526 
527 if (pending_close) {
528 close();
529 }
530}
531 
532template <typename CompressionContext>
533void ZlibUtil::CompressionStream<CompressionContext>::close() {
534 pending_close = writing;
535 if (writing) {
536 return;
537 }
538 closed = true;
539 JSG_ASSERT(initialized, Error, "Closing before initialized"_kj);
540 // Context is closed on the destructor of the CompressionContext.
541}
542 
543template <typename CompressionContext>
544bool ZlibUtil::CompressionStream<CompressionContext>::checkError(jsg::Lock& js) {
545 KJ_IF_SOME(error, context()->getError()) {
546 emitError(js, kj::mv(error));
547 return false;
548 }
549 return true;
550}
551 
552template <typename CompressionContext>
553void ZlibUtil::CompressionStream<CompressionContext>::initializeStream(
554 jsg::Lock& js, jsg::JsArrayBufferView& _writeResult, jsg::Function<void()> _writeCallback) {
555 writeResult = _writeResult.addRef(js);
556 writeCallback = kj::mv(_writeCallback);
557 initialized = true;
558}
559 
560template <typename CompressionContext>
561void ZlibUtil::CompressionStream<CompressionContext>::updateWriteResult(jsg::Lock& js) {
562 KJ_IF_SOME(wr, writeResult) {
563 auto result = wr.getHandle(js);
564 auto ptr = result.template asArrayPtr<uint32_t>();
565 JSG_REQUIRE(ptr.size() >= 2, Error, "Invalid write result buffer"_kj);
566 context()->getAfterWriteResult(&ptr[1], &ptr[0]);
567 }
568}
569 
570template <typename CompressionContext>
571template <bool async>
572void ZlibUtil::CompressionStream<CompressionContext>::write(jsg::Lock& js,
573 int flush,
574 jsg::Optional<kj::Array<kj::byte>> input,
575 uint32_t inputOffset,
576 uint32_t inputLength,
577 kj::Array<kj::byte> output,
578 uint32_t outputOffset,
579 uint32_t outputLength) {
580 if (flush != Z_NO_FLUSH && flush != Z_PARTIAL_FLUSH && flush != Z_SYNC_FLUSH &&
581 flush != Z_FULL_FLUSH && flush != Z_FINISH && flush != Z_BLOCK) {
582 JSG_FAIL_REQUIRE(Error, "Invalid flush value");
583 }
584 
585 // Use default values if input is not determined
586 if (input == kj::none) {
587 inputLength = 0;
588 inputOffset = 0;
589 }
590 
591 auto input_ensured = input.map([](auto& val) { return val.asPtr(); }).orDefault({});
592 
593 // Check for integer overflow...
594 JSG_REQUIRE(inputOffset + inputLength >= inputOffset, Error, "Input access it not within bounds");
595 JSG_REQUIRE(
596 outputOffset + outputLength >= outputOffset, Error, "Input access it not within bounds");
597 JSG_REQUIRE(IsWithinBounds(inputOffset, inputLength, input_ensured.size()), Error,
598 "Input access is not within bounds"_kj);
599 JSG_REQUIRE(IsWithinBounds(outputOffset, outputLength, output.size()), Error,
600 "Output access is not within bounds"_kj);
601 
602 writeStream<async>(js, flush, input_ensured.slice(inputOffset, inputOffset + inputLength),
603 output.slice(outputOffset, outputOffset + outputLength));
604}
605 
606template <typename CompressionContext>
607void ZlibUtil::CompressionStream<CompressionContext>::reset(jsg::Lock& js) {
608 KJ_IF_SOME(error, context()->resetStream()) {
609 emitError(js, kj::mv(error));
610 }
611}
612 
613jsg::Ref<ZlibUtil::ZlibStream> ZlibUtil::ZlibStream::constructor(
614 jsg::Lock& js, ZlibModeValue mode) {
615 return js.alloc<ZlibStream>(static_cast<ZlibMode>(mode), js.getExternalMemoryTarget());
616}
617 
618void ZlibUtil::ZlibStream::initialize(jsg::Lock& js,
619 int windowBits,
620 int level,
621 int memLevel,
622 int strategy,
623 jsg::JsArrayBufferView writeState,
624 jsg::Function<void()> writeCallback,
625 jsg::Optional<kj::Array<kj::byte>> dictionary) {
626 initializeStream(js, writeState, kj::mv(writeCallback));
627 allocator.configure(context()->getStream());
628 context()->initialize(level, windowBits, memLevel, strategy, kj::mv(dictionary));
629}
630 
631void ZlibUtil::ZlibStream::params(jsg::Lock& js, int _level, int _strategy) {
632 context()->setParams(_level, _strategy);
633 KJ_IF_SOME(err, context()->getError()) {
634 emitError(js, kj::mv(err));
635 }
636}
637 
638void BrotliContext::setBuffers(kj::ArrayPtr<kj::byte> input, kj::ArrayPtr<kj::byte> output) {
639 nextIn = reinterpret_cast<const uint8_t*>(input.begin());
640 nextOut = output.begin();
641 availIn = input.size();
642 availOut = output.size();
643}
644 
645void BrotliContext::setInputBuffer(kj::ArrayPtr<const kj::byte> input) {
646 nextIn = input.begin();
647 availIn = input.size();
648}
649 
650void BrotliContext::setOutputBuffer(kj::ArrayPtr<kj::byte> output) {
651 nextOut = output.begin();
652 availOut = output.size();
653}
654 
655uint BrotliContext::getAvailOut() const {
656 return availOut;
657}
658 
659void BrotliContext::setFlush(int _flush) {
660 flush = static_cast<BrotliEncoderOperation>(_flush);
661}
662 
663void BrotliContext::getAfterWriteResult(uint32_t* _availIn, uint32_t* _availOut) const {
664 *_availIn = availIn;
665 *_availOut = availOut;
666}
667 
668BrotliEncoderContext::BrotliEncoderContext(ZlibMode _mode): BrotliContext(_mode) {
669 auto instance = BrotliEncoderCreateInstance(alloc_brotli, free_brotli, alloc_opaque_brotli);
670 state = kj::disposeWith<BrotliEncoderDestroyInstance>(instance);
671}
672 
673void BrotliEncoderContext::work() {
674 JSG_REQUIRE(mode == ZlibMode::BROTLI_ENCODE, Error, "Mode should be BROTLI_ENCODE"_kj);
675 JSG_REQUIRE_NONNULL(state.get(), Error, "State should not be empty"_kj);
676 
677 const uint8_t* internalNext = nextIn;
678 lastResult = BrotliEncoderCompressStream(
679 state.get(), flush, &availIn, &internalNext, &availOut, &nextOut, nullptr);
680 nextIn += internalNext - nextIn;
681 
682 streamEnd = lastResult && BrotliEncoderIsFinished(state.get());
683}
684 
685kj::Maybe<CompressionError> BrotliEncoderContext::initialize(
686 brotli_alloc_func init_alloc_func, brotli_free_func init_free_func, void* init_opaque_func) {
687 alloc_brotli = init_alloc_func;
688 free_brotli = init_free_func;
689 alloc_opaque_brotli = init_opaque_func;
690 
691 auto instance = BrotliEncoderCreateInstance(alloc_brotli, free_brotli, alloc_opaque_brotli);
692 state = kj::disposeWith<BrotliEncoderDestroyInstance>(kj::mv(instance));
693 
694 if (state.get() == nullptr) {
695 return CompressionError(
696 "Could not initialize Brotli instance"_kj, "ERR_ZLIB_INITIALIZATION_FAILED"_kj, -1);
697 }
698 
699 return kj::none;
700}
701 
702kj::Maybe<CompressionError> BrotliEncoderContext::resetStream() {
703 return initialize(alloc_brotli, free_brotli, alloc_opaque_brotli);
704}
705 
706kj::Maybe<CompressionError> BrotliEncoderContext::setParams(int key, uint32_t value) {
707 if (!BrotliEncoderSetParameter(state.get(), static_cast<BrotliEncoderParameter>(key), value)) {
708 return CompressionError("Setting parameter failed", "ERR_BROTLI_PARAM_SET_FAILED", -1);
709 }
710 
711 return kj::none;
712}
713 
714kj::Maybe<CompressionError> BrotliEncoderContext::getError() const {
715 if (!lastResult) {
716 return CompressionError("Compression failed", "ERR_BROTLI_COMPRESSION_FAILED", -1);
717 }
718 
719 return kj::none;
720}
721 
722bool BrotliEncoderContext::isStreamEnd() const {
723 return streamEnd;
724}
725 
726BrotliDecoderContext::BrotliDecoderContext(ZlibMode _mode): BrotliContext(_mode) {
727 auto instance = BrotliDecoderCreateInstance(alloc_brotli, free_brotli, alloc_opaque_brotli);
728 state = kj::disposeWith<BrotliDecoderDestroyInstance>(instance);
729}
730 
731kj::Maybe<CompressionError> BrotliDecoderContext::initialize(
732 brotli_alloc_func init_alloc_func, brotli_free_func init_free_func, void* init_opaque_func) {
733 alloc_brotli = init_alloc_func;
734 free_brotli = init_free_func;
735 alloc_opaque_brotli = init_opaque_func;
736 
737 auto instance = BrotliDecoderCreateInstance(alloc_brotli, free_brotli, alloc_opaque_brotli);
738 state = kj::disposeWith<BrotliDecoderDestroyInstance>(kj::mv(instance));
739 
740 if (state.get() == nullptr) {
741 return CompressionError(
742 "Could not initialize Brotli instance", "ERR_ZLIB_INITIALIZATION_FAILED", -1);
743 }
744 
745 return kj::none;
746}
747 
748void BrotliDecoderContext::work() {
749 JSG_REQUIRE(mode == ZlibMode::BROTLI_DECODE, Error, "Mode should have been BROTLI_DECODE"_kj);
750 JSG_REQUIRE_NONNULL(state.get(), Error, "State should not be empty"_kj);
751 const uint8_t* internalNext = nextIn;
752 lastResult = BrotliDecoderDecompressStream(
753 state.get(), &availIn, &internalNext, &availOut, &nextOut, nullptr);
754 nextIn += internalNext - nextIn;
755 
756 if (lastResult == BROTLI_DECODER_RESULT_ERROR) {
757 error = BrotliDecoderGetErrorCode(state.get());
758 errorString = kj::str("ERR_", BrotliDecoderErrorString(error));
759 }
760}
761 
762kj::Maybe<CompressionError> BrotliDecoderContext::resetStream() {
763 return initialize(alloc_brotli, free_brotli, alloc_opaque_brotli);
764}
765 
766kj::Maybe<CompressionError> BrotliDecoderContext::setParams(int key, uint32_t value) {
767 if (!BrotliDecoderSetParameter(state.get(), static_cast<BrotliDecoderParameter>(key), value)) {
768 return CompressionError("Setting parameter failed", "ERR_BROTLI_PARAM_SET_FAILED", -1);
769 }
770 
771 return kj::none;
772}
773 
774kj::Maybe<CompressionError> BrotliDecoderContext::getError() const {
775 if (error != BROTLI_DECODER_NO_ERROR) {
776 return CompressionError("Compression failed", errorString, -1);
777 }
778 
779 if (flush == BROTLI_OPERATION_FINISH && lastResult == BROTLI_DECODER_RESULT_NEEDS_MORE_INPUT) {
780 // Match zlib behavior, as brotli doesn't have its own code for this.
781 return CompressionError("Unexpected end of file", "Z_BUF_ERROR", Z_BUF_ERROR);
782 }
783 
784 return kj::none;
785}
786 
787bool BrotliDecoderContext::isStreamEnd() const {
788 return lastResult == BROTLI_DECODER_RESULT_SUCCESS;
789}
790 
791// =======================================================================================
792// Zstd Implementation
793 
794void ZstdContext::setBuffers(kj::ArrayPtr<kj::byte> input, kj::ArrayPtr<kj::byte> output) {
795 setInputBuffer(input);
796 setOutputBuffer(output);
797}
798 
799void ZstdContext::setInputBuffer(kj::ArrayPtr<const kj::byte> input) {
800 input_.src = input.begin();
801 input_.size = input.size();
802 input_.pos = 0;
803}
804 
805void ZstdContext::setOutputBuffer(kj::ArrayPtr<kj::byte> output) {
806 output_.dst = output.begin();
807 output_.size = output.size();
808 output_.pos = 0;
809}
810 
811void ZstdContext::setFlush(int flush) {
812 KJ_DASSERT(flush >= ZSTD_e_continue && flush <= ZSTD_e_end,
813 "flush must be a valid ZSTD_EndDirective value");
814 flush_ = static_cast<ZSTD_EndDirective>(flush);
815}
816 
817kj::uint ZstdContext::getAvailOut() const {
818 return output_.size - output_.pos;
819}
820 
821void ZstdContext::getAfterWriteResult(uint32_t* availIn, uint32_t* availOut) const {
822 *availIn = input_.size - input_.pos;
823 *availOut = output_.size - output_.pos;
824}
825 
826namespace {
827// Helper to check ZSTD errors and return a CompressionError if present.
828// Also sets the error code in the provided reference for later retrieval.
829kj::Maybe<CompressionError> zstdCheckError(
830 size_t result, ZSTD_ErrorCode& error, kj::StringPtr errorCode) {
831 if (ZSTD_isError(result)) {
832 error = ZSTD_getErrorCode(result);
833 return CompressionError(ZSTD_getErrorName(result), errorCode, -1);
834 }
835 return kj::none;
836}
837 
838// Wrappers for ZSTD free functions that return void (for use with kj::disposeWith).
839void zstdFreeCCtx(ZSTD_CCtx* cctx) {
840 ZSTD_freeCCtx(cctx);
841}
842void zstdFreeDCtx(ZSTD_DCtx* dctx) {
843 ZSTD_freeDCtx(dctx);
844}
845} // namespace
846 
847ZstdEncoderContext::ZstdEncoderContext(ZlibMode _mode)
848 : ZstdContext(_mode),
849 cctx_(kj::disposeWith<zstdFreeCCtx>(ZSTD_createCCtx())) {}
850 
851kj::Maybe<CompressionError> ZstdEncoderContext::initialize(uint64_t pledgedSrcSize) {
852 if (cctx_.get() == nullptr) {
853 return CompressionError(
854 "Could not initialize Zstd instance"_kj, "ERR_ZLIB_INITIALIZATION_FAILED"_kj, -1);
855 }
856 
857 if (pledgedSrcSize != ZSTD_CONTENTSIZE_UNKNOWN) {
858 size_t result = ZSTD_CCtx_setPledgedSrcSize(cctx_.get(), pledgedSrcSize);
859 KJ_IF_SOME(err, zstdCheckError(result, error_, "ERR_ZSTD_COMPRESSION_FAILED"_kj)) {
860 return kj::mv(err);
861 }
862 }
863 
864 return kj::none;
865}
866 
867void ZstdEncoderContext::work() {
868 JSG_REQUIRE(mode == ZlibMode::ZSTD_ENCODE, Error, "Mode should be ZSTD_ENCODE"_kj);
869 JSG_REQUIRE(cctx_.get() != nullptr, Error, "Zstd context should not be null"_kj);
870 
871 lastResult = ZSTD_compressStream2(cctx_.get(), &output_, &input_, flush_);
872 
873 if (ZSTD_isError(lastResult)) {
874 error_ = ZSTD_getErrorCode(lastResult);
875 }
876}
877 
878kj::Maybe<CompressionError> ZstdEncoderContext::resetStream() {
879 if (cctx_.get() != nullptr) {
880 size_t result = ZSTD_CCtx_reset(cctx_.get(), ZSTD_reset_session_only);
881 KJ_IF_SOME(err, zstdCheckError(result, error_, "ERR_ZSTD_COMPRESSION_FAILED"_kj)) {
882 return kj::mv(err);
883 }
884 }
885 return kj::none;
886}
887 
888kj::Maybe<CompressionError> ZstdEncoderContext::setParams(int key, int value) {
889 KJ_DASSERT(key >= ZSTD_c_compressionLevel,
890 "key must be a valid ZSTD_cParameter (first valid value is ZSTD_c_compressionLevel)");
891 size_t result = ZSTD_CCtx_setParameter(cctx_.get(), static_cast<ZSTD_cParameter>(key), value);
892 if (ZSTD_isError(result)) {
893 return CompressionError(kj::str("Setting parameter failed: ", ZSTD_getErrorName(result)),
894 "ERR_ZSTD_PARAM_SET_FAILED"_kj, -1);
895 }
896 return kj::none;
897}
898 
899kj::Maybe<CompressionError> ZstdEncoderContext::getError() const {
900 if (error_ != ZSTD_error_no_error) {
901 return CompressionError(kj::str("Zstd compression failed: ", ZSTD_getErrorString(error_)),
902 kj::str("ERR_ZSTD_COMPRESSION_FAILED"), -1);
903 }
904 
905 if (flush_ == ZSTD_e_end && lastResult != 0) {
906 // lastResult > 0 means more output is needed, which shouldn't happen at end
907 return CompressionError("Unexpected end of file"_kj, "Z_BUF_ERROR"_kj, Z_BUF_ERROR);
908 }
909 
910 return kj::none;
911}
912 
913bool ZstdEncoderContext::isStreamEnd() const {
914 // ZSTD_compressStream2 returns 0 when flush_ == ZSTD_e_end and the frame is fully flushed.
915 return !ZSTD_isError(lastResult) && lastResult == 0;
916}
917 
918ZstdDecoderContext::ZstdDecoderContext(ZlibMode _mode)
919 : ZstdContext(_mode),
920 dctx_(kj::disposeWith<zstdFreeDCtx>(ZSTD_createDCtx())) {}
921 
922kj::Maybe<CompressionError> ZstdDecoderContext::initialize() {
923 // dctx_ is created in the constructor. It can only be nullptr if ZSTD_createDCtx()
924 // failed due to memory allocation failure.
925 if (dctx_.get() == nullptr) {
926 return CompressionError(
927 "Could not initialize Zstd instance"_kj, "ERR_ZLIB_INITIALIZATION_FAILED"_kj, -1);
928 }
929 
930 return kj::none;
931}
932 
933void ZstdDecoderContext::work() {
934 JSG_REQUIRE(mode == ZlibMode::ZSTD_DECODE, Error, "Mode should be ZSTD_DECODE"_kj);
935 JSG_REQUIRE(dctx_.get() != nullptr, Error, "Zstd context should not be null"_kj);
936 
937 lastResult = ZSTD_decompressStream(dctx_.get(), &output_, &input_);
938 
939 if (ZSTD_isError(lastResult)) {
940 error_ = ZSTD_getErrorCode(lastResult);
941 } else if (input_.size > 0) {
942 // Track whether we're mid-frame: lastResult > 0 means more data needed,
943 // lastResult == 0 means frame is complete.
944 frameInProgress_ = (lastResult > 0);
945 }
946}
947 
948kj::Maybe<CompressionError> ZstdDecoderContext::resetStream() {
949 if (dctx_.get() != nullptr) {
950 size_t result = ZSTD_DCtx_reset(dctx_.get(), ZSTD_reset_session_only);
951 KJ_IF_SOME(err, zstdCheckError(result, error_, "ERR_ZSTD_DECOMPRESSION_FAILED"_kj)) {
952 return kj::mv(err);
953 }
954 }
955 frameInProgress_ = false;
956 return kj::none;
957}
958 
959kj::Maybe<CompressionError> ZstdDecoderContext::setParams(int key, int value) {
960 KJ_DASSERT(dctx_.get() != nullptr, "Zstd decompression context should not be null");
961 size_t result = ZSTD_DCtx_setParameter(dctx_.get(), static_cast<ZSTD_dParameter>(key), value);
962 if (ZSTD_isError(result)) {
963 return CompressionError(kj::str("Setting parameter failed: ", ZSTD_getErrorName(result)),
964 "ERR_ZSTD_PARAM_SET_FAILED"_kj, -1);
965 }
966 return kj::none;
967}
968 
969kj::Maybe<CompressionError> ZstdDecoderContext::getError() const {
970 if (error_ != ZSTD_error_no_error) {
971 return CompressionError(kj::str("Zstd decompression failed: ", ZSTD_getErrorString(error_)),
972 kj::str("ERR_ZSTD_DECOMPRESSION_FAILED"), -1);
973 }
974 
975 // If this is the final flush, we're mid-frame (frame was started but never
976 // completed), and the output buffer is not full (decoder had space but
977 // couldn't produce more output), the input was truncated.
978 if (flush_ == ZSTD_e_end && frameInProgress_ && output_.pos < output_.size) {
979 return CompressionError("unexpected end of file"_kj, "ERR_ZSTD_DECOMPRESSION_FAILED"_kj, -1);
980 }
981 
982 return kj::none;
983}
984 
985bool ZstdDecoderContext::isStreamEnd() const {
986 // ZSTD_decompressStream returns 0 when a frame is completely decoded and fully flushed.
987 return !ZSTD_isError(lastResult) && lastResult == 0;
988}
989 
990template <typename CompressionContext>
991jsg::Ref<ZlibUtil::ZstdCompressionStream<CompressionContext>> ZlibUtil::ZstdCompressionStream<
992 CompressionContext>::constructor(jsg::Lock& js, ZlibModeValue mode) {
993 return js.alloc<ZstdCompressionStream>(static_cast<ZlibMode>(mode), js.getExternalMemoryTarget());
994}
995 
996template <typename CompressionContext>
997bool ZlibUtil::ZstdCompressionStream<CompressionContext>::initialize(jsg::Lock& js,
998 jsg::JsArrayBufferView params,
999 jsg::JsArrayBufferView writeResult,
1000 jsg::Function<void()> writeCallback,
1001 jsg::Optional<uint64_t> pledgedSrcSize) {
1002 this->initializeStream(js, writeResult, kj::mv(writeCallback));
1003 
1004 uint64_t srcSize = pledgedSrcSize.orDefault(ZSTD_CONTENTSIZE_UNKNOWN);
1005 
1006 kj::Maybe<CompressionError> maybeError;
1007 if constexpr (CompressionContext::Mode == ZlibMode::ZSTD_ENCODE) {
1008 maybeError = this->context()->initialize(srcSize);
1009 } else {
1010 maybeError = this->context()->initialize();
1011 }
1012 
1013 KJ_IF_SOME(err, maybeError) {
1014 this->emitError(js, kj::mv(err));
1015 return false;
1016 }
1017 
1018 auto results = params.template asArrayPtr<int>();
1019 
1020 for (size_t i = 0; i < results.size(); i++) {
1021 if (results[i] == -1) {
1022 continue;
1023 }
1024 
1025 KJ_IF_SOME(err, this->context()->setParams(i, results[i])) {
1026 this->emitError(js, kj::mv(err));
1027 return false;
1028 }
1029 }
1030 return true;
1031}
1032 
1033template <typename CompressionContext>
1034jsg::Ref<ZlibUtil::BrotliCompressionStream<CompressionContext>> ZlibUtil::BrotliCompressionStream<
1035 CompressionContext>::constructor(jsg::Lock& js, ZlibModeValue mode) {
1036 return js.alloc<BrotliCompressionStream>(
1037 static_cast<ZlibMode>(mode), js.getExternalMemoryTarget());
1038}
1039 
1040template <typename CompressionContext>
1041bool ZlibUtil::BrotliCompressionStream<CompressionContext>::initialize(jsg::Lock& js,
1042 jsg::JsArrayBufferView params,
1043 jsg::JsArrayBufferView writeResult,
1044 jsg::Function<void()> writeCallback) {
1045 this->initializeStream(js, writeResult, kj::mv(writeCallback));
1046 auto maybeError = this->context()->initialize(
1047 CompressionAllocator::AllocForBrotli, CompressionAllocator::FreeForZlib, &this->allocator);
1048 
1049 KJ_IF_SOME(err, maybeError) {
1050 this->emitError(js, kj::mv(err));
1051 return false;
1052 }
1053 
1054 auto results = params.template asArrayPtr<uint32_t>();
1055 
1056 for (int i = 0; i < results.size(); i++) {
1057 if (results[i] == static_cast<uint32_t>(-1)) {
1058 continue;
1059 }
1060 
1061 KJ_IF_SOME(err, this->context()->setParams(i, results[i])) {
1062 this->emitError(js, kj::mv(err));
1063 return false;
1064 }
1065 }
1066 return true;
1067}
1068 
1069namespace {
1070template <typename Context>
1071static kj::Array<kj::byte> syncProcessBuffer(Context& ctx, GrowableBuffer& result) {
1072 do {
1073 result.addChunk();
1074 ctx.setOutputBuffer(kj::ArrayPtr(result.end(), result.available()));
1075 
1076 ctx.work();
1077 
1078 KJ_IF_SOME(error, ctx.getError()) {
1079 JSG_FAIL_REQUIRE(Error, error.message);
1080 }
1081 
1082 result.adjustUnused(ctx.getAvailOut());
1083 
1084 if (ctx.getAvailOut() == 0 && result.atMaxCapacity()) {
1085 // The output buffer was completely filled and has reached maxOutputLength.
1086 // If the stream is done, the output just happened to fit exactly — return it.
1087 // Otherwise the decompressed data exceeds maxOutputLength.
1088 JSG_REQUIRE(ctx.isStreamEnd(), RangeError, "Memory limit exceeded");
1089 break;
1090 }
1091 } while (ctx.getAvailOut() == 0);
1092 
1093 return result.releaseAsArray();
1094}
1095} // namespace
1096 
1097kj::Array<kj::byte> ZlibUtil::zlibSync(
1098 jsg::Lock& js, ZlibUtil::InputSource data, ZlibContext::Options opts, ZlibModeValue mode) {
1099 // Any use of zlib APIs constitutes an implicit dependency on Allocator which must
1100 // remain alive until the zlib stream is destroyed
1101 CompressionAllocator allocator(js.getExternalMemoryTarget());
1102 ZlibContext ctx(static_cast<ZlibMode>(mode));
1103 allocator.configure(ctx.getStream());
1104 
1105 auto chunkSize = opts.chunkSize.orDefault(ZLIB_PERFORMANT_CHUNK_SIZE);
1106 auto maxOutputLength = opts.maxOutputLength.orDefault(Z_MAX_CHUNK);
1107 
1108 // TODO(soon): Extend JSG_REQUIRE so we can pass the full level of info NodeJS provides, like the code field
1109 JSG_REQUIRE(Z_MIN_CHUNK <= chunkSize && chunkSize <= Z_MAX_CHUNK, RangeError,
1110 kj::str("The value of \"options.chunkSize\" is out of range. It must be >= ", Z_MIN_CHUNK,
1111 " and <= ", Z_MAX_CHUNK, ". Received ", chunkSize));
1112 JSG_REQUIRE(maxOutputLength >= 1 && maxOutputLength <= Z_MAX_CHUNK, RangeError,
1113 kj::str("The value of \"options.maxOutputLength\" is out of range. It must be >= 1 and <= ",
1114 Z_MAX_CHUNK, ". Received ", maxOutputLength));
1115 GrowableBuffer result(ZLIB_PERFORMANT_CHUNK_SIZE, maxOutputLength);
1116 
1117 ctx.initialize(opts.level.orDefault(Z_DEFAULT_LEVEL),
1118 opts.windowBits.orDefault(Z_DEFAULT_WINDOWBITS), opts.memLevel.orDefault(Z_DEFAULT_MEMLEVEL),
1119 opts.strategy.orDefault(Z_DEFAULT_STRATEGY), kj::mv(opts.dictionary));
1120 
1121 auto flush = opts.flush.orDefault(Z_NO_FLUSH);
1122 JSG_REQUIRE(Z_NO_FLUSH <= flush && flush <= Z_TREES, RangeError,
1123 kj::str("The value of \"options.flush\" is out of range. It must be >= ", Z_NO_FLUSH,
1124 " and <= ", Z_TREES, ". Received ", flush));
1125 
1126 auto finishFlush = opts.finishFlush.orDefault(Z_FINISH);
1127 JSG_REQUIRE(Z_NO_FLUSH <= finishFlush && finishFlush <= Z_TREES, RangeError,
1128 kj::str("The value of \"options.finishFlush\" is out of range. It must be >= ", Z_NO_FLUSH,
1129 " and <= ", Z_TREES, ". Received ", flush));
1130 
1131 ctx.setFlush(finishFlush);
1132 ctx.setInputBuffer(getInputFromSource(data));
1133 return syncProcessBuffer(ctx, result);
1134}
1135 
1136void ZlibUtil::zlibWithCallback(jsg::Lock& js,
1137 InputSource data,
1138 ZlibContext::Options options,
1139 ZlibModeValue mode,
1140 CompressCallback cb) {
1141 // Capture only relevant errors so they can be passed to the callback
1142 auto res = js.tryCatch([&]() {
1143 return CompressCallbackArg(zlibSync(js, kj::mv(data), kj::mv(options), mode));
1144 }, [&](jsg::Value&& exception) {
1145 return CompressCallbackArg(jsg::JsValue(exception.getHandle(js)));
1146 });
1147 
1148 // Ensure callback is invoked only once
1149 cb(js, kj::mv(res));
1150}
1151 
1152template <typename Context>
1153kj::Array<kj::byte> ZlibUtil::brotliSync(
1154 jsg::Lock& js, InputSource data, BrotliContext::Options opts) {
1155 // Any use of brotli APIs constitutes an implicit dependency on Allocator which must
1156 // remain alive until the brotli state is destroyed
1157 CompressionAllocator allocator(js.getExternalMemoryTarget());
1158 Context ctx(Context::Mode);
1159 
1160 auto chunkSize = opts.chunkSize.orDefault(ZLIB_PERFORMANT_CHUNK_SIZE);
1161 auto maxOutputLength = opts.maxOutputLength.orDefault(Z_MAX_CHUNK);
1162 
1163 // TODO(soon): Extend JSG_REQUIRE so we can pass the full level of info NodeJS provides, like the code field
1164 JSG_REQUIRE(Z_MIN_CHUNK <= chunkSize && chunkSize <= Z_MAX_CHUNK, RangeError,
1165 kj::str("The value of \"options.chunkSize\" is out of range. It must be >= ", Z_MIN_CHUNK,
1166 " and <= ", Z_MAX_CHUNK, ". Received ", chunkSize));
1167 JSG_REQUIRE(maxOutputLength >= 1 && maxOutputLength <= Z_MAX_CHUNK, RangeError,
1168 kj::str("The value of \"options.maxOutputLength\" is out of range. It must be >= 1 and <= ",
1169 Z_MAX_CHUNK, ". Received ", maxOutputLength));
1170 GrowableBuffer result(ZLIB_PERFORMANT_CHUNK_SIZE, maxOutputLength);
1171 
1172 KJ_IF_SOME(err,
1173 ctx.initialize(
1174 CompressionAllocator::AllocForBrotli, CompressionAllocator::FreeForZlib, &allocator)) {
1175 JSG_FAIL_REQUIRE(Error, err.message);
1176 }
1177 
1178 KJ_IF_SOME(params, opts.params) {
1179 for (const auto& field: params.fields) {
1180 KJ_IF_SOME(err, ctx.setParams(field.name.parseAs<int>(), field.value)) {
1181 JSG_FAIL_REQUIRE(Error, err.message);
1182 }
1183 }
1184 }
1185 
1186 auto flush = opts.flush.orDefault(BROTLI_OPERATION_PROCESS);
1187 JSG_REQUIRE(BROTLI_OPERATION_PROCESS <= flush && flush <= BROTLI_OPERATION_EMIT_METADATA,
1188 RangeError,
1189 kj::str("The value of \"options.flush\" is out of range. It must be >= ",
1190 BROTLI_OPERATION_PROCESS, " and <= ", BROTLI_OPERATION_EMIT_METADATA, ". Received ",
1191 flush));
1192 
1193 auto finishFlush = opts.finishFlush.orDefault(BROTLI_OPERATION_FINISH);
1194 JSG_REQUIRE(
1195 BROTLI_OPERATION_PROCESS <= finishFlush && finishFlush <= BROTLI_OPERATION_EMIT_METADATA,
1196 RangeError,
1197 kj::str("The value of \"options.finishFlush\" is out of range. It must be >= ",
1198 BROTLI_OPERATION_PROCESS, " and <= ", BROTLI_OPERATION_EMIT_METADATA, ". Received ",
1199 finishFlush));
1200 
1201 ctx.setFlush(finishFlush);
1202 ctx.setInputBuffer(getInputFromSource(data));
1203 return syncProcessBuffer(ctx, result);
1204}
1205 
1206template <typename Context>
1207void ZlibUtil::brotliWithCallback(
1208 jsg::Lock& js, InputSource data, BrotliContext::Options options, CompressCallback cb) {
1209 // Capture only relevant errors so they can be passed to the callback
1210 auto res = js.tryCatch([&]() {
1211 return CompressCallbackArg(brotliSync<Context>(js, kj::mv(data), kj::mv(options)));
1212 }, [&](jsg::Value&& exception) {
1213 return CompressCallbackArg(jsg::JsValue(exception.getHandle(js)));
1214 });
1215 
1216 // Ensure callback is invoked only once
1217 cb(js, kj::mv(res));
1218}
1219 
1220template <typename Context>
1221kj::Array<kj::byte> ZlibUtil::zstdSync(jsg::Lock& js, InputSource data, ZstdContext::Options opts) {
1222 Context ctx(Context::Mode);
1223 
1224 auto chunkSize = opts.chunkSize.orDefault(ZLIB_PERFORMANT_CHUNK_SIZE);
1225 auto maxOutputLength = opts.maxOutputLength.orDefault(Z_MAX_CHUNK);
1226 
1227 JSG_REQUIRE(Z_MIN_CHUNK <= chunkSize && chunkSize <= Z_MAX_CHUNK, RangeError,
1228 kj::str("The value of \"options.chunkSize\" is out of range. It must be >= ", Z_MIN_CHUNK,
1229 " and <= ", Z_MAX_CHUNK, ". Received ", chunkSize));
1230 JSG_REQUIRE(maxOutputLength >= 1 && maxOutputLength <= Z_MAX_CHUNK, RangeError,
1231 kj::str("The value of \"options.maxOutputLength\" is out of range. It must be >= 1 and <= ",
1232 Z_MAX_CHUNK, ". Received ", maxOutputLength));
1233 GrowableBuffer result(ZLIB_PERFORMANT_CHUNK_SIZE, maxOutputLength);
1234 
1235 // Initialize the context
1236 if constexpr (Context::Mode == ZlibMode::ZSTD_ENCODE) {
1237 uint64_t pledgedSrcSize = opts.pledgedSrcSize.orDefault(ZSTD_CONTENTSIZE_UNKNOWN);
1238 KJ_IF_SOME(err, ctx.initialize(pledgedSrcSize)) {
1239 JSG_FAIL_REQUIRE(Error, err.message);
1240 }
1241 } else {
1242 KJ_IF_SOME(err, ctx.initialize()) {
1243 JSG_FAIL_REQUIRE(Error, err.message);
1244 }
1245 }
1246 
1247 // Set parameters
1248 KJ_IF_SOME(params, opts.params) {
1249 for (const auto& field: params.fields) {
1250 KJ_IF_SOME(err, ctx.setParams(field.name.parseAs<int>(), field.value)) {
1251 JSG_FAIL_REQUIRE(Error, err.message);
1252 }
1253 }
1254 }
1255 
1256 auto flush = opts.flush.orDefault(ZSTD_e_continue);
1257 JSG_REQUIRE(ZSTD_e_continue <= flush && flush <= ZSTD_e_end, RangeError,
1258 kj::str("The value of \"options.flush\" is out of range. It must be >= ", ZSTD_e_continue,
1259 " and <= ", ZSTD_e_end, ". Received ", flush));
1260 
1261 auto finishFlush = opts.finishFlush.orDefault(ZSTD_e_end);
1262 JSG_REQUIRE(ZSTD_e_continue <= finishFlush && finishFlush <= ZSTD_e_end, RangeError,
1263 kj::str("The value of \"options.finishFlush\" is out of range. It must be >= ",
1264 ZSTD_e_continue, " and <= ", ZSTD_e_end, ". Received ", finishFlush));
1265 
1266 ctx.setFlush(finishFlush);
1267 ctx.setInputBuffer(getInputFromSource(data));
1268 return syncProcessBuffer(ctx, result);
1269}
1270 
1271template <typename Context>
1272void ZlibUtil::zstdWithCallback(
1273 jsg::Lock& js, InputSource data, ZstdContext::Options options, CompressCallback cb) {
1274 // Capture only relevant errors so they can be passed to the callback
1275 auto res = js.tryCatch([&]() {
1276 return CompressCallbackArg(zstdSync<Context>(js, kj::mv(data), kj::mv(options)));
1277 }, [&](jsg::Value&& exception) {
1278 return CompressCallbackArg(jsg::JsValue(exception.getHandle(js)));
1279 });
1280 
1281 // Ensure callback is invoked only once
1282 cb(js, kj::mv(res));
1283}
1284 
1285#ifndef CREATE_TEMPLATE
1286#define CREATE_TEMPLATE(T) \
1287 template class ZlibUtil::CompressionStream<T>; \
1288 template void ZlibUtil::CompressionStream<T>::write<false>(jsg::Lock & js, int flush, \
1289 jsg::Optional<kj::Array<kj::byte>> input, uint32_t inputOffset, uint32_t inputLength, \
1290 kj::Array<kj::byte> output, uint32_t outputOffset, uint32_t outputLength); \
1291 template void ZlibUtil::CompressionStream<T>::write<true>(jsg::Lock & js, int flush, \
1292 jsg::Optional<kj::Array<kj::byte>> input, uint32_t inputOffset, uint32_t inputLength, \
1293 kj::Array<kj::byte> output, uint32_t outputOffset, uint32_t outputLength);
1294 
1295CREATE_TEMPLATE(ZlibContext)
1296CREATE_TEMPLATE(BrotliEncoderContext)
1297CREATE_TEMPLATE(BrotliDecoderContext)
1298CREATE_TEMPLATE(ZstdEncoderContext)
1299CREATE_TEMPLATE(ZstdDecoderContext)
1300 
1301template class ZlibUtil::BrotliCompressionStream<BrotliEncoderContext>;
1302template class ZlibUtil::BrotliCompressionStream<BrotliDecoderContext>;
1303 
1304template class ZlibUtil::ZstdCompressionStream<ZstdEncoderContext>;
1305template class ZlibUtil::ZstdCompressionStream<ZstdDecoderContext>;
1306 
1307template kj::Array<kj::byte> ZlibUtil::brotliSync<BrotliEncoderContext>(
1308 jsg::Lock& js, InputSource data, BrotliContext::Options opts);
1309template kj::Array<kj::byte> ZlibUtil::brotliSync<BrotliDecoderContext>(
1310 jsg::Lock& js, InputSource data, BrotliContext::Options opts);
1311template void ZlibUtil::brotliWithCallback<BrotliEncoderContext>(
1312 jsg::Lock& js, InputSource data, BrotliContext::Options options, CompressCallback cb);
1313template void ZlibUtil::brotliWithCallback<BrotliDecoderContext>(
1314 jsg::Lock& js, InputSource data, BrotliContext::Options options, CompressCallback cb);
1315 
1316template kj::Array<kj::byte> ZlibUtil::zstdSync<ZstdEncoderContext>(
1317 jsg::Lock& js, InputSource data, ZstdContext::Options opts);
1318template kj::Array<kj::byte> ZlibUtil::zstdSync<ZstdDecoderContext>(
1319 jsg::Lock& js, InputSource data, ZstdContext::Options opts);
1320template void ZlibUtil::zstdWithCallback<ZstdEncoderContext>(
1321 jsg::Lock& js, InputSource data, ZstdContext::Options options, CompressCallback cb);
1322template void ZlibUtil::zstdWithCallback<ZstdDecoderContext>(
1323 jsg::Lock& js, InputSource data, ZstdContext::Options options, CompressCallback cb);
1324#undef CREATE_TEMPLATE
1325#endif
1326} // namespace workerd::api::node