// Copyright (c) 2017-2022 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 // Copyright Joyent and Node contributors. All rights reserved. MIT license. #include "zlib-util.h" #include "util.h" // The following implementation is adapted from Node.js // and therefore follows Node.js style as opposed to kj style. // Latest implementation of Node.js zlib can be found at: // https://github.com/nodejs/node/blob/main/src/node_zlib.cc namespace workerd::api::node { kj::ArrayPtr getInputFromSource(const ZlibUtil::InputSource& data) { KJ_SWITCH_ONEOF(data) { KJ_CASE_ONEOF(dataBuf, kj::Array) { JSG_REQUIRE(dataBuf.size() < Z_MAX_CHUNK, RangeError, "Memory limit exceeded"_kj); return dataBuf.asPtr(); } KJ_CASE_ONEOF(dataStr, jsg::NonCoercible) { JSG_REQUIRE(dataStr.value.size() < Z_MAX_CHUNK, RangeError, "Memory limit exceeded"_kj); return dataStr.value.asBytes(); } } KJ_UNREACHABLE; } uint32_t ZlibUtil::crc32Sync(InputSource data, uint32_t value) { auto dataPtr = getInputFromSource(data); return crc32(value, dataPtr.begin(), dataPtr.size()); } namespace { class GrowableBuffer final { // A copy of kj::Vector with some additional methods for use as a growable buffer with a maximum // size public: inline explicit GrowableBuffer(size_t _chunkSize, size_t _maxCapacity) : maxCapacity(_maxCapacity) { auto maxChunkSize = kj::min(_chunkSize, maxCapacity); builder = kj::heapArrayBuilder(maxChunkSize); chunkSize = maxChunkSize; } size_t size() const { return builder.size(); } bool empty() const { return size() == 0; } size_t capacity() const { return builder.capacity(); } size_t available() const { return capacity() - size(); } kj::byte* begin() KJ_LIFETIMEBOUND { return builder.begin(); } kj::byte* end() KJ_LIFETIMEBOUND { return builder.end(); } kj::Array releaseAsArray() { // TODO(perf): Avoid a copy/move by allowing Array to point to incomplete space? if (!builder.isFull()) { setCapacity(size()); } return builder.finish(); } void adjustUnused(size_t unused) { resize(capacity() - unused); } void resize(size_t size) { if (size > builder.capacity()) grow(size); builder.resize(size); } void addChunk() { reserve(size() + chunkSize); } void reserve(size_t size) { if (size > builder.capacity()) { grow(size); } } bool atMaxCapacity() const { return size() >= maxCapacity; } private: kj::ArrayBuilder builder; size_t chunkSize; size_t maxCapacity; void grow(size_t minCapacity = 0) { JSG_REQUIRE(minCapacity <= maxCapacity, RangeError, "Memory limit exceeded"); setCapacity(kj::min(maxCapacity, kj::max(minCapacity, capacity() == 0 ? 4 : capacity() * 2))); } void setCapacity(size_t newSize) { if (builder.size() > newSize) { builder.truncate(newSize); } kj::ArrayBuilder newBuilder = kj::heapArrayBuilder(newSize); newBuilder.addAll(kj::mv(builder)); builder = kj::mv(newBuilder); } }; } // namespace void ZlibContext::initialize(int _level, int _windowBits, int _memLevel, int _strategy, jsg::Optional> _dictionary) { if (!((_windowBits == 0) && (mode == ZlibMode::INFLATE || mode == ZlibMode::GUNZIP || mode == ZlibMode::UNZIP))) { JSG_ASSERT(_windowBits >= Z_MIN_WINDOWBITS && _windowBits <= Z_MAX_WINDOWBITS, RangeError, kj::str("The value of \"options.windowBits\" is out of range. It must be >= ", Z_MIN_WINDOWBITS, " and <= ", Z_MAX_WINDOWBITS, ". Received ", _windowBits)); } JSG_REQUIRE(_level >= Z_MIN_LEVEL && _level <= Z_MAX_LEVEL, RangeError, kj::str("The value of \"options.level\" is out of range. It must be >= ", Z_MIN_LEVEL, " and <= ", Z_MAX_LEVEL, ". Received ", _level)); JSG_REQUIRE(_memLevel >= Z_MIN_MEMLEVEL && _memLevel <= Z_MAX_MEMLEVEL, RangeError, kj::str("The value of \"options.memLevel\" is out of range. It must be >= ", Z_MIN_MEMLEVEL, " and <= ", Z_MAX_MEMLEVEL, ". Received ", _memLevel)); JSG_REQUIRE(_strategy == Z_FILTERED || _strategy == Z_HUFFMAN_ONLY || _strategy == Z_RLE || _strategy == Z_FIXED || _strategy == Z_DEFAULT_STRATEGY, Error, "invalid strategy"_kj); level = _level; windowBits = _windowBits; memLevel = _memLevel; strategy = _strategy; flush = Z_NO_FLUSH; err = Z_OK; switch (mode) { case ZlibMode::GZIP: case ZlibMode::GUNZIP: windowBits += 16; break; case ZlibMode::UNZIP: windowBits += 32; break; case ZlibMode::DEFLATERAW: case ZlibMode::INFLATERAW: windowBits *= -1; break; default: break; } KJ_IF_SOME(dict, _dictionary) { dictionary = kj::mv(dict); } } kj::Maybe ZlibContext::getError() const { // Acceptable error states depend on the type of zlib stream. switch (err) { case Z_OK: case Z_BUF_ERROR: if (stream.avail_out != 0 && flush == Z_FINISH) { return constructError("unexpected end of file"_kj); } break; case Z_STREAM_END: // normal statuses, not fatal break; case Z_NEED_DICT: if (dictionary.empty()) { return constructError("Missing dictionary"_kj); } else { return constructError("Bad dictionary"_kj); } default: // something else. return constructError("Zlib error"); } return {}; } kj::Maybe ZlibContext::setDictionary() { if (dictionary.empty()) { return kj::none; } err = Z_OK; switch (mode) { case ZlibMode::DEFLATE: case ZlibMode::DEFLATERAW: err = deflateSetDictionary(&stream, dictionary.begin(), dictionary.size()); break; case ZlibMode::INFLATERAW: err = inflateSetDictionary(&stream, dictionary.begin(), dictionary.size()); break; default: break; } if (err != Z_OK) { return constructError("Failed to set dictionary"_kj); } return kj::none; } bool ZlibContext::initializeZlib() { if (initialized) { return false; } switch (mode) { case ZlibMode::DEFLATE: case ZlibMode::GZIP: case ZlibMode::DEFLATERAW: err = deflateInit2(&stream, level, Z_DEFLATED, windowBits, memLevel, strategy); break; case ZlibMode::INFLATE: case ZlibMode::GUNZIP: case ZlibMode::INFLATERAW: case ZlibMode::UNZIP: err = inflateInit2(&stream, windowBits); break; default: KJ_UNREACHABLE; } if (err != Z_OK) { dictionary.clear(); mode = ZlibMode::NONE; return true; } setDictionary(); initialized = true; return true; } kj::Maybe ZlibContext::resetStream() { bool initialized_now = initializeZlib(); if (initialized_now && err != Z_OK) { return constructError("Failed to init stream before reset"); } err = Z_OK; switch (mode) { case ZlibMode::DEFLATE: case ZlibMode::DEFLATERAW: case ZlibMode::GZIP: err = deflateReset(&stream); break; case ZlibMode::INFLATE: case ZlibMode::INFLATERAW: case ZlibMode::GUNZIP: err = inflateReset(&stream); break; default: break; } if (err != Z_OK) { return constructError("Failed to reset stream"_kj); } return setDictionary(); } void ZlibContext::work() { bool initialized_now = initializeZlib(); if (initialized_now && err != Z_OK) { return; } const Bytef* next_expected_header_byte = nullptr; // If the avail_out is left at 0, then it means that it ran out // of room. If there was avail_out left over, then it means // that all the input was consumed. switch (mode) { case ZlibMode::DEFLATE: case ZlibMode::GZIP: case ZlibMode::DEFLATERAW: err = deflate(&stream, flush); break; case ZlibMode::UNZIP: if (stream.avail_in > 0) { next_expected_header_byte = stream.next_in; } switch (gzip_id_bytes_read) { case 0: if (next_expected_header_byte == nullptr) { break; } if (*next_expected_header_byte == GZIP_HEADER_ID1) { gzip_id_bytes_read = 1; next_expected_header_byte++; if (stream.avail_in == 1) { // The only available byte was already read. break; } } else { mode = ZlibMode::INFLATE; break; } [[fallthrough]]; case 1: if (next_expected_header_byte == nullptr) { break; } if (*next_expected_header_byte == GZIP_HEADER_ID2) { gzip_id_bytes_read = 2; mode = ZlibMode::GUNZIP; } else { // There is no actual difference between INFLATE and INFLATERAW // (after initialization). mode = ZlibMode::INFLATE; } break; default: JSG_FAIL_REQUIRE(Error, "Invalid number of gzip magic number bytes read"); } [[fallthrough]]; case ZlibMode::INFLATE: case ZlibMode::GUNZIP: case ZlibMode::INFLATERAW: err = inflate(&stream, flush); // If data was encoded with dictionary (INFLATERAW will have it set in // SetDictionary, don't repeat that here) if (mode != ZlibMode::INFLATERAW && err == Z_NEED_DICT && !dictionary.empty()) { // Load it err = inflateSetDictionary(&stream, dictionary.begin(), dictionary.size()); if (err == Z_OK) { // And try to decode again err = inflate(&stream, flush); } else if (err == Z_DATA_ERROR) { // Both inflateSetDictionary() and inflate() return Z_DATA_ERROR. // Make it possible for After() to tell a bad dictionary from bad // input. err = Z_NEED_DICT; } } while (stream.avail_in > 0 && mode == ZlibMode::GUNZIP && err == Z_STREAM_END && stream.next_in[0] != 0x00) { // Bytes remain in input buffer. Perhaps this is another compressed // member in the same archive, or just trailing garbage. // Trailing zero bytes are okay, though, since they are frequently // used for padding. resetStream(); err = inflate(&stream, flush); } break; default: KJ_UNREACHABLE; } } kj::Maybe ZlibContext::setParams(int _level, int _strategy) { bool initialized_now = initializeZlib(); if (initialized_now && err != Z_OK) { return constructError("Failed to init stream before set parameters"); } err = Z_OK; switch (mode) { case ZlibMode::DEFLATE: case ZlibMode::DEFLATERAW: err = deflateParams(&stream, _level, _strategy); break; default: break; } if (err != Z_OK && err != Z_BUF_ERROR) { return constructError("Failed to set parameters"); } return kj::none; } ZlibContext::~ZlibContext() noexcept(false) { if (!initialized) { return; } auto status = Z_OK; switch (mode) { case ZlibMode::DEFLATE: case ZlibMode::DEFLATERAW: case ZlibMode::GZIP: status = deflateEnd(&stream); break; case ZlibMode::INFLATE: case ZlibMode::INFLATERAW: case ZlibMode::GUNZIP: case ZlibMode::UNZIP: status = inflateEnd(&stream); break; default: break; } JSG_REQUIRE( status == Z_OK || status == Z_DATA_ERROR, Error, "Uncaught error on closing zlib stream"); } void ZlibContext::setBuffers(kj::ArrayPtr input, kj::ArrayPtr output) { stream.avail_in = input.size(); stream.next_in = input.begin(); stream.avail_out = output.size(); stream.next_out = output.begin(); } void ZlibContext::setInputBuffer(kj::ArrayPtr input) { // The define Z_CONST is not set, so zlib always takes mutable pointers stream.next_in = const_cast(input.begin()); stream.avail_in = input.size(); } void ZlibContext::setOutputBuffer(kj::ArrayPtr output) { stream.next_out = output.begin(); stream.avail_out = output.size(); } template jsg::Ref> ZlibUtil::CompressionStream< CompressionContext>::constructor(jsg::Lock& js, ZlibModeValue mode) { return js.alloc(static_cast(mode), js.getExternalMemoryTarget()); } template ZlibUtil::CompressionStream::~CompressionStream() { JSG_ASSERT(!writing, Error, "Writing to compression stream"_kj); close(); } template void ZlibUtil::CompressionStream::emitError( jsg::Lock& js, const CompressionError& error) { KJ_IF_SOME(onError, errorHandler) { onError(js, error.err, kj::mv(error.code), kj::mv(error.message)); } writing = false; if (pending_close) { close(); } } template template void ZlibUtil::CompressionStream::writeStream( jsg::Lock& js, int flush, kj::ArrayPtr input, kj::ArrayPtr output) { JSG_REQUIRE(initialized, Error, "Writing before initializing"_kj); JSG_REQUIRE(!closed, Error, "Already finalized"_kj); JSG_REQUIRE(!writing, Error, "Writing is in progress"_kj); JSG_REQUIRE(!pending_close, Error, "Pending close"_kj); writing = true; context()->setBuffers(input, output); context()->setFlush(flush); // Clear buffer pointers from the compression context when this scope exits. // The input/output kj::Array parameters are backed by V8 BackingStores // whose lifetimes are tied to their JavaScript ArrayBuffer objects. Once // this method returns, the kj::Array destructor releases its shared_ptr // to the BackingStore. If the JS buffer subsequently becomes unreachable // and is garbage-collected, the BackingStore is freed — leaving any // retained pointers (e.g. z_stream.next_out) dangling. A later call to // deflateParams() (via params()) could then write to freed memory. // // Using KJ_DEFER ensures the pointers are cleared even if an exception // is thrown (e.g. from updateWriteResult() when the writeState buffer is // too small). Without this, the exception would unwind the stack before // reaching an explicit clearBuffers() call, leaving stale pointers. KJ_DEFER(context()->clearBuffers()); if constexpr (!async) { context()->work(); if (checkError(js)) { writing = false; updateWriteResult(js); } return; } // On Node.js, this is called as a result of `ScheduleWork()` call. // Since, we implement the whole thing as sync, we're going to ahead and call the whole thing here. context()->work(); // This is implemented slightly differently in Node.js // Node.js calls AfterThreadPoolWork(). // Ref: https://github.com/nodejs/node/blob/9edf4a0856681a7665bd9dcf2ca7cac252784b98/src/node_zlib.cc#L402 writing = false; if (!checkError(js)) { return; } updateWriteResult(js); KJ_IF_SOME(cb, writeCallback) { cb(js); } if (pending_close) { close(); } } template void ZlibUtil::CompressionStream::close() { pending_close = writing; if (writing) { return; } closed = true; JSG_ASSERT(initialized, Error, "Closing before initialized"_kj); // Context is closed on the destructor of the CompressionContext. } template bool ZlibUtil::CompressionStream::checkError(jsg::Lock& js) { KJ_IF_SOME(error, context()->getError()) { emitError(js, kj::mv(error)); return false; } return true; } template void ZlibUtil::CompressionStream::initializeStream( jsg::Lock& js, jsg::JsArrayBufferView& _writeResult, jsg::Function _writeCallback) { writeResult = _writeResult.addRef(js); writeCallback = kj::mv(_writeCallback); initialized = true; } template void ZlibUtil::CompressionStream::updateWriteResult(jsg::Lock& js) { KJ_IF_SOME(wr, writeResult) { auto result = wr.getHandle(js); auto ptr = result.template asArrayPtr(); JSG_REQUIRE(ptr.size() >= 2, Error, "Invalid write result buffer"_kj); context()->getAfterWriteResult(&ptr[1], &ptr[0]); } } template template void ZlibUtil::CompressionStream::write(jsg::Lock& js, int flush, jsg::Optional> input, uint32_t inputOffset, uint32_t inputLength, kj::Array output, uint32_t outputOffset, uint32_t outputLength) { if (flush != Z_NO_FLUSH && flush != Z_PARTIAL_FLUSH && flush != Z_SYNC_FLUSH && flush != Z_FULL_FLUSH && flush != Z_FINISH && flush != Z_BLOCK) { JSG_FAIL_REQUIRE(Error, "Invalid flush value"); } // Use default values if input is not determined if (input == kj::none) { inputLength = 0; inputOffset = 0; } auto input_ensured = input.map([](auto& val) { return val.asPtr(); }).orDefault({}); // Check for integer overflow... JSG_REQUIRE(inputOffset + inputLength >= inputOffset, Error, "Input access it not within bounds"); JSG_REQUIRE( outputOffset + outputLength >= outputOffset, Error, "Input access it not within bounds"); JSG_REQUIRE(IsWithinBounds(inputOffset, inputLength, input_ensured.size()), Error, "Input access is not within bounds"_kj); JSG_REQUIRE(IsWithinBounds(outputOffset, outputLength, output.size()), Error, "Output access is not within bounds"_kj); writeStream(js, flush, input_ensured.slice(inputOffset, inputOffset + inputLength), output.slice(outputOffset, outputOffset + outputLength)); } template void ZlibUtil::CompressionStream::reset(jsg::Lock& js) { KJ_IF_SOME(error, context()->resetStream()) { emitError(js, kj::mv(error)); } } jsg::Ref ZlibUtil::ZlibStream::constructor( jsg::Lock& js, ZlibModeValue mode) { return js.alloc(static_cast(mode), js.getExternalMemoryTarget()); } void ZlibUtil::ZlibStream::initialize(jsg::Lock& js, int windowBits, int level, int memLevel, int strategy, jsg::JsArrayBufferView writeState, jsg::Function writeCallback, jsg::Optional> dictionary) { initializeStream(js, writeState, kj::mv(writeCallback)); allocator.configure(context()->getStream()); context()->initialize(level, windowBits, memLevel, strategy, kj::mv(dictionary)); } void ZlibUtil::ZlibStream::params(jsg::Lock& js, int _level, int _strategy) { context()->setParams(_level, _strategy); KJ_IF_SOME(err, context()->getError()) { emitError(js, kj::mv(err)); } } void BrotliContext::setBuffers(kj::ArrayPtr input, kj::ArrayPtr output) { nextIn = reinterpret_cast(input.begin()); nextOut = output.begin(); availIn = input.size(); availOut = output.size(); } void BrotliContext::setInputBuffer(kj::ArrayPtr input) { nextIn = input.begin(); availIn = input.size(); } void BrotliContext::setOutputBuffer(kj::ArrayPtr output) { nextOut = output.begin(); availOut = output.size(); } uint BrotliContext::getAvailOut() const { return availOut; } void BrotliContext::setFlush(int _flush) { flush = static_cast(_flush); } void BrotliContext::getAfterWriteResult(uint32_t* _availIn, uint32_t* _availOut) const { *_availIn = availIn; *_availOut = availOut; } BrotliEncoderContext::BrotliEncoderContext(ZlibMode _mode): BrotliContext(_mode) { auto instance = BrotliEncoderCreateInstance(alloc_brotli, free_brotli, alloc_opaque_brotli); state = kj::disposeWith(instance); } void BrotliEncoderContext::work() { JSG_REQUIRE(mode == ZlibMode::BROTLI_ENCODE, Error, "Mode should be BROTLI_ENCODE"_kj); JSG_REQUIRE_NONNULL(state.get(), Error, "State should not be empty"_kj); const uint8_t* internalNext = nextIn; lastResult = BrotliEncoderCompressStream( state.get(), flush, &availIn, &internalNext, &availOut, &nextOut, nullptr); nextIn += internalNext - nextIn; streamEnd = lastResult && BrotliEncoderIsFinished(state.get()); } kj::Maybe BrotliEncoderContext::initialize( brotli_alloc_func init_alloc_func, brotli_free_func init_free_func, void* init_opaque_func) { alloc_brotli = init_alloc_func; free_brotli = init_free_func; alloc_opaque_brotli = init_opaque_func; auto instance = BrotliEncoderCreateInstance(alloc_brotli, free_brotli, alloc_opaque_brotli); state = kj::disposeWith(kj::mv(instance)); if (state.get() == nullptr) { return CompressionError( "Could not initialize Brotli instance"_kj, "ERR_ZLIB_INITIALIZATION_FAILED"_kj, -1); } return kj::none; } kj::Maybe BrotliEncoderContext::resetStream() { return initialize(alloc_brotli, free_brotli, alloc_opaque_brotli); } kj::Maybe BrotliEncoderContext::setParams(int key, uint32_t value) { if (!BrotliEncoderSetParameter(state.get(), static_cast(key), value)) { return CompressionError("Setting parameter failed", "ERR_BROTLI_PARAM_SET_FAILED", -1); } return kj::none; } kj::Maybe BrotliEncoderContext::getError() const { if (!lastResult) { return CompressionError("Compression failed", "ERR_BROTLI_COMPRESSION_FAILED", -1); } return kj::none; } bool BrotliEncoderContext::isStreamEnd() const { return streamEnd; } BrotliDecoderContext::BrotliDecoderContext(ZlibMode _mode): BrotliContext(_mode) { auto instance = BrotliDecoderCreateInstance(alloc_brotli, free_brotli, alloc_opaque_brotli); state = kj::disposeWith(instance); } kj::Maybe BrotliDecoderContext::initialize( brotli_alloc_func init_alloc_func, brotli_free_func init_free_func, void* init_opaque_func) { alloc_brotli = init_alloc_func; free_brotli = init_free_func; alloc_opaque_brotli = init_opaque_func; auto instance = BrotliDecoderCreateInstance(alloc_brotli, free_brotli, alloc_opaque_brotli); state = kj::disposeWith(kj::mv(instance)); if (state.get() == nullptr) { return CompressionError( "Could not initialize Brotli instance", "ERR_ZLIB_INITIALIZATION_FAILED", -1); } return kj::none; } void BrotliDecoderContext::work() { JSG_REQUIRE(mode == ZlibMode::BROTLI_DECODE, Error, "Mode should have been BROTLI_DECODE"_kj); JSG_REQUIRE_NONNULL(state.get(), Error, "State should not be empty"_kj); const uint8_t* internalNext = nextIn; lastResult = BrotliDecoderDecompressStream( state.get(), &availIn, &internalNext, &availOut, &nextOut, nullptr); nextIn += internalNext - nextIn; if (lastResult == BROTLI_DECODER_RESULT_ERROR) { error = BrotliDecoderGetErrorCode(state.get()); errorString = kj::str("ERR_", BrotliDecoderErrorString(error)); } } kj::Maybe BrotliDecoderContext::resetStream() { return initialize(alloc_brotli, free_brotli, alloc_opaque_brotli); } kj::Maybe BrotliDecoderContext::setParams(int key, uint32_t value) { if (!BrotliDecoderSetParameter(state.get(), static_cast(key), value)) { return CompressionError("Setting parameter failed", "ERR_BROTLI_PARAM_SET_FAILED", -1); } return kj::none; } kj::Maybe BrotliDecoderContext::getError() const { if (error != BROTLI_DECODER_NO_ERROR) { return CompressionError("Compression failed", errorString, -1); } if (flush == BROTLI_OPERATION_FINISH && lastResult == BROTLI_DECODER_RESULT_NEEDS_MORE_INPUT) { // Match zlib behavior, as brotli doesn't have its own code for this. return CompressionError("Unexpected end of file", "Z_BUF_ERROR", Z_BUF_ERROR); } return kj::none; } bool BrotliDecoderContext::isStreamEnd() const { return lastResult == BROTLI_DECODER_RESULT_SUCCESS; } // ======================================================================================= // Zstd Implementation void ZstdContext::setBuffers(kj::ArrayPtr input, kj::ArrayPtr output) { setInputBuffer(input); setOutputBuffer(output); } void ZstdContext::setInputBuffer(kj::ArrayPtr input) { input_.src = input.begin(); input_.size = input.size(); input_.pos = 0; } void ZstdContext::setOutputBuffer(kj::ArrayPtr output) { output_.dst = output.begin(); output_.size = output.size(); output_.pos = 0; } void ZstdContext::setFlush(int flush) { KJ_DASSERT(flush >= ZSTD_e_continue && flush <= ZSTD_e_end, "flush must be a valid ZSTD_EndDirective value"); flush_ = static_cast(flush); } kj::uint ZstdContext::getAvailOut() const { return output_.size - output_.pos; } void ZstdContext::getAfterWriteResult(uint32_t* availIn, uint32_t* availOut) const { *availIn = input_.size - input_.pos; *availOut = output_.size - output_.pos; } namespace { // Helper to check ZSTD errors and return a CompressionError if present. // Also sets the error code in the provided reference for later retrieval. kj::Maybe zstdCheckError( size_t result, ZSTD_ErrorCode& error, kj::StringPtr errorCode) { if (ZSTD_isError(result)) { error = ZSTD_getErrorCode(result); return CompressionError(ZSTD_getErrorName(result), errorCode, -1); } return kj::none; } // Wrappers for ZSTD free functions that return void (for use with kj::disposeWith). void zstdFreeCCtx(ZSTD_CCtx* cctx) { ZSTD_freeCCtx(cctx); } void zstdFreeDCtx(ZSTD_DCtx* dctx) { ZSTD_freeDCtx(dctx); } } // namespace ZstdEncoderContext::ZstdEncoderContext(ZlibMode _mode) : ZstdContext(_mode), cctx_(kj::disposeWith(ZSTD_createCCtx())) {} kj::Maybe ZstdEncoderContext::initialize(uint64_t pledgedSrcSize) { if (cctx_.get() == nullptr) { return CompressionError( "Could not initialize Zstd instance"_kj, "ERR_ZLIB_INITIALIZATION_FAILED"_kj, -1); } if (pledgedSrcSize != ZSTD_CONTENTSIZE_UNKNOWN) { size_t result = ZSTD_CCtx_setPledgedSrcSize(cctx_.get(), pledgedSrcSize); KJ_IF_SOME(err, zstdCheckError(result, error_, "ERR_ZSTD_COMPRESSION_FAILED"_kj)) { return kj::mv(err); } } return kj::none; } void ZstdEncoderContext::work() { JSG_REQUIRE(mode == ZlibMode::ZSTD_ENCODE, Error, "Mode should be ZSTD_ENCODE"_kj); JSG_REQUIRE(cctx_.get() != nullptr, Error, "Zstd context should not be null"_kj); lastResult = ZSTD_compressStream2(cctx_.get(), &output_, &input_, flush_); if (ZSTD_isError(lastResult)) { error_ = ZSTD_getErrorCode(lastResult); } } kj::Maybe ZstdEncoderContext::resetStream() { if (cctx_.get() != nullptr) { size_t result = ZSTD_CCtx_reset(cctx_.get(), ZSTD_reset_session_only); KJ_IF_SOME(err, zstdCheckError(result, error_, "ERR_ZSTD_COMPRESSION_FAILED"_kj)) { return kj::mv(err); } } return kj::none; } kj::Maybe ZstdEncoderContext::setParams(int key, int value) { KJ_DASSERT(key >= ZSTD_c_compressionLevel, "key must be a valid ZSTD_cParameter (first valid value is ZSTD_c_compressionLevel)"); size_t result = ZSTD_CCtx_setParameter(cctx_.get(), static_cast(key), value); if (ZSTD_isError(result)) { return CompressionError(kj::str("Setting parameter failed: ", ZSTD_getErrorName(result)), "ERR_ZSTD_PARAM_SET_FAILED"_kj, -1); } return kj::none; } kj::Maybe ZstdEncoderContext::getError() const { if (error_ != ZSTD_error_no_error) { return CompressionError(kj::str("Zstd compression failed: ", ZSTD_getErrorString(error_)), kj::str("ERR_ZSTD_COMPRESSION_FAILED"), -1); } if (flush_ == ZSTD_e_end && lastResult != 0) { // lastResult > 0 means more output is needed, which shouldn't happen at end return CompressionError("Unexpected end of file"_kj, "Z_BUF_ERROR"_kj, Z_BUF_ERROR); } return kj::none; } bool ZstdEncoderContext::isStreamEnd() const { // ZSTD_compressStream2 returns 0 when flush_ == ZSTD_e_end and the frame is fully flushed. return !ZSTD_isError(lastResult) && lastResult == 0; } ZstdDecoderContext::ZstdDecoderContext(ZlibMode _mode) : ZstdContext(_mode), dctx_(kj::disposeWith(ZSTD_createDCtx())) {} kj::Maybe ZstdDecoderContext::initialize() { // dctx_ is created in the constructor. It can only be nullptr if ZSTD_createDCtx() // failed due to memory allocation failure. if (dctx_.get() == nullptr) { return CompressionError( "Could not initialize Zstd instance"_kj, "ERR_ZLIB_INITIALIZATION_FAILED"_kj, -1); } return kj::none; } void ZstdDecoderContext::work() { JSG_REQUIRE(mode == ZlibMode::ZSTD_DECODE, Error, "Mode should be ZSTD_DECODE"_kj); JSG_REQUIRE(dctx_.get() != nullptr, Error, "Zstd context should not be null"_kj); lastResult = ZSTD_decompressStream(dctx_.get(), &output_, &input_); if (ZSTD_isError(lastResult)) { error_ = ZSTD_getErrorCode(lastResult); } else if (input_.size > 0) { // Track whether we're mid-frame: lastResult > 0 means more data needed, // lastResult == 0 means frame is complete. frameInProgress_ = (lastResult > 0); } } kj::Maybe ZstdDecoderContext::resetStream() { if (dctx_.get() != nullptr) { size_t result = ZSTD_DCtx_reset(dctx_.get(), ZSTD_reset_session_only); KJ_IF_SOME(err, zstdCheckError(result, error_, "ERR_ZSTD_DECOMPRESSION_FAILED"_kj)) { return kj::mv(err); } } frameInProgress_ = false; return kj::none; } kj::Maybe ZstdDecoderContext::setParams(int key, int value) { KJ_DASSERT(dctx_.get() != nullptr, "Zstd decompression context should not be null"); size_t result = ZSTD_DCtx_setParameter(dctx_.get(), static_cast(key), value); if (ZSTD_isError(result)) { return CompressionError(kj::str("Setting parameter failed: ", ZSTD_getErrorName(result)), "ERR_ZSTD_PARAM_SET_FAILED"_kj, -1); } return kj::none; } kj::Maybe ZstdDecoderContext::getError() const { if (error_ != ZSTD_error_no_error) { return CompressionError(kj::str("Zstd decompression failed: ", ZSTD_getErrorString(error_)), kj::str("ERR_ZSTD_DECOMPRESSION_FAILED"), -1); } // If this is the final flush, we're mid-frame (frame was started but never // completed), and the output buffer is not full (decoder had space but // couldn't produce more output), the input was truncated. if (flush_ == ZSTD_e_end && frameInProgress_ && output_.pos < output_.size) { return CompressionError("unexpected end of file"_kj, "ERR_ZSTD_DECOMPRESSION_FAILED"_kj, -1); } return kj::none; } bool ZstdDecoderContext::isStreamEnd() const { // ZSTD_decompressStream returns 0 when a frame is completely decoded and fully flushed. return !ZSTD_isError(lastResult) && lastResult == 0; } template jsg::Ref> ZlibUtil::ZstdCompressionStream< CompressionContext>::constructor(jsg::Lock& js, ZlibModeValue mode) { return js.alloc(static_cast(mode), js.getExternalMemoryTarget()); } template bool ZlibUtil::ZstdCompressionStream::initialize(jsg::Lock& js, jsg::JsArrayBufferView params, jsg::JsArrayBufferView writeResult, jsg::Function writeCallback, jsg::Optional pledgedSrcSize) { this->initializeStream(js, writeResult, kj::mv(writeCallback)); uint64_t srcSize = pledgedSrcSize.orDefault(ZSTD_CONTENTSIZE_UNKNOWN); kj::Maybe maybeError; if constexpr (CompressionContext::Mode == ZlibMode::ZSTD_ENCODE) { maybeError = this->context()->initialize(srcSize); } else { maybeError = this->context()->initialize(); } KJ_IF_SOME(err, maybeError) { this->emitError(js, kj::mv(err)); return false; } auto results = params.template asArrayPtr(); for (size_t i = 0; i < results.size(); i++) { if (results[i] == -1) { continue; } KJ_IF_SOME(err, this->context()->setParams(i, results[i])) { this->emitError(js, kj::mv(err)); return false; } } return true; } template jsg::Ref> ZlibUtil::BrotliCompressionStream< CompressionContext>::constructor(jsg::Lock& js, ZlibModeValue mode) { return js.alloc( static_cast(mode), js.getExternalMemoryTarget()); } template bool ZlibUtil::BrotliCompressionStream::initialize(jsg::Lock& js, jsg::JsArrayBufferView params, jsg::JsArrayBufferView writeResult, jsg::Function writeCallback) { this->initializeStream(js, writeResult, kj::mv(writeCallback)); auto maybeError = this->context()->initialize( CompressionAllocator::AllocForBrotli, CompressionAllocator::FreeForZlib, &this->allocator); KJ_IF_SOME(err, maybeError) { this->emitError(js, kj::mv(err)); return false; } auto results = params.template asArrayPtr(); for (int i = 0; i < results.size(); i++) { if (results[i] == static_cast(-1)) { continue; } KJ_IF_SOME(err, this->context()->setParams(i, results[i])) { this->emitError(js, kj::mv(err)); return false; } } return true; } namespace { template static kj::Array syncProcessBuffer(Context& ctx, GrowableBuffer& result) { do { result.addChunk(); ctx.setOutputBuffer(kj::ArrayPtr(result.end(), result.available())); ctx.work(); KJ_IF_SOME(error, ctx.getError()) { JSG_FAIL_REQUIRE(Error, error.message); } result.adjustUnused(ctx.getAvailOut()); if (ctx.getAvailOut() == 0 && result.atMaxCapacity()) { // The output buffer was completely filled and has reached maxOutputLength. // If the stream is done, the output just happened to fit exactly — return it. // Otherwise the decompressed data exceeds maxOutputLength. JSG_REQUIRE(ctx.isStreamEnd(), RangeError, "Memory limit exceeded"); break; } } while (ctx.getAvailOut() == 0); return result.releaseAsArray(); } } // namespace kj::Array ZlibUtil::zlibSync( jsg::Lock& js, ZlibUtil::InputSource data, ZlibContext::Options opts, ZlibModeValue mode) { // Any use of zlib APIs constitutes an implicit dependency on Allocator which must // remain alive until the zlib stream is destroyed CompressionAllocator allocator(js.getExternalMemoryTarget()); ZlibContext ctx(static_cast(mode)); allocator.configure(ctx.getStream()); auto chunkSize = opts.chunkSize.orDefault(ZLIB_PERFORMANT_CHUNK_SIZE); auto maxOutputLength = opts.maxOutputLength.orDefault(Z_MAX_CHUNK); // TODO(soon): Extend JSG_REQUIRE so we can pass the full level of info NodeJS provides, like the code field JSG_REQUIRE(Z_MIN_CHUNK <= chunkSize && chunkSize <= Z_MAX_CHUNK, RangeError, kj::str("The value of \"options.chunkSize\" is out of range. It must be >= ", Z_MIN_CHUNK, " and <= ", Z_MAX_CHUNK, ". Received ", chunkSize)); JSG_REQUIRE(maxOutputLength >= 1 && maxOutputLength <= Z_MAX_CHUNK, RangeError, kj::str("The value of \"options.maxOutputLength\" is out of range. It must be >= 1 and <= ", Z_MAX_CHUNK, ". Received ", maxOutputLength)); GrowableBuffer result(ZLIB_PERFORMANT_CHUNK_SIZE, maxOutputLength); ctx.initialize(opts.level.orDefault(Z_DEFAULT_LEVEL), opts.windowBits.orDefault(Z_DEFAULT_WINDOWBITS), opts.memLevel.orDefault(Z_DEFAULT_MEMLEVEL), opts.strategy.orDefault(Z_DEFAULT_STRATEGY), kj::mv(opts.dictionary)); auto flush = opts.flush.orDefault(Z_NO_FLUSH); JSG_REQUIRE(Z_NO_FLUSH <= flush && flush <= Z_TREES, RangeError, kj::str("The value of \"options.flush\" is out of range. It must be >= ", Z_NO_FLUSH, " and <= ", Z_TREES, ". Received ", flush)); auto finishFlush = opts.finishFlush.orDefault(Z_FINISH); JSG_REQUIRE(Z_NO_FLUSH <= finishFlush && finishFlush <= Z_TREES, RangeError, kj::str("The value of \"options.finishFlush\" is out of range. It must be >= ", Z_NO_FLUSH, " and <= ", Z_TREES, ". Received ", flush)); ctx.setFlush(finishFlush); ctx.setInputBuffer(getInputFromSource(data)); return syncProcessBuffer(ctx, result); } void ZlibUtil::zlibWithCallback(jsg::Lock& js, InputSource data, ZlibContext::Options options, ZlibModeValue mode, CompressCallback cb) { // Capture only relevant errors so they can be passed to the callback auto res = js.tryCatch([&]() { return CompressCallbackArg(zlibSync(js, kj::mv(data), kj::mv(options), mode)); }, [&](jsg::Value&& exception) { return CompressCallbackArg(jsg::JsValue(exception.getHandle(js))); }); // Ensure callback is invoked only once cb(js, kj::mv(res)); } template kj::Array ZlibUtil::brotliSync( jsg::Lock& js, InputSource data, BrotliContext::Options opts) { // Any use of brotli APIs constitutes an implicit dependency on Allocator which must // remain alive until the brotli state is destroyed CompressionAllocator allocator(js.getExternalMemoryTarget()); Context ctx(Context::Mode); auto chunkSize = opts.chunkSize.orDefault(ZLIB_PERFORMANT_CHUNK_SIZE); auto maxOutputLength = opts.maxOutputLength.orDefault(Z_MAX_CHUNK); // TODO(soon): Extend JSG_REQUIRE so we can pass the full level of info NodeJS provides, like the code field JSG_REQUIRE(Z_MIN_CHUNK <= chunkSize && chunkSize <= Z_MAX_CHUNK, RangeError, kj::str("The value of \"options.chunkSize\" is out of range. It must be >= ", Z_MIN_CHUNK, " and <= ", Z_MAX_CHUNK, ". Received ", chunkSize)); JSG_REQUIRE(maxOutputLength >= 1 && maxOutputLength <= Z_MAX_CHUNK, RangeError, kj::str("The value of \"options.maxOutputLength\" is out of range. It must be >= 1 and <= ", Z_MAX_CHUNK, ". Received ", maxOutputLength)); GrowableBuffer result(ZLIB_PERFORMANT_CHUNK_SIZE, maxOutputLength); KJ_IF_SOME(err, ctx.initialize( CompressionAllocator::AllocForBrotli, CompressionAllocator::FreeForZlib, &allocator)) { JSG_FAIL_REQUIRE(Error, err.message); } KJ_IF_SOME(params, opts.params) { for (const auto& field: params.fields) { KJ_IF_SOME(err, ctx.setParams(field.name.parseAs(), field.value)) { JSG_FAIL_REQUIRE(Error, err.message); } } } auto flush = opts.flush.orDefault(BROTLI_OPERATION_PROCESS); JSG_REQUIRE(BROTLI_OPERATION_PROCESS <= flush && flush <= BROTLI_OPERATION_EMIT_METADATA, RangeError, kj::str("The value of \"options.flush\" is out of range. It must be >= ", BROTLI_OPERATION_PROCESS, " and <= ", BROTLI_OPERATION_EMIT_METADATA, ". Received ", flush)); auto finishFlush = opts.finishFlush.orDefault(BROTLI_OPERATION_FINISH); JSG_REQUIRE( BROTLI_OPERATION_PROCESS <= finishFlush && finishFlush <= BROTLI_OPERATION_EMIT_METADATA, RangeError, kj::str("The value of \"options.finishFlush\" is out of range. It must be >= ", BROTLI_OPERATION_PROCESS, " and <= ", BROTLI_OPERATION_EMIT_METADATA, ". Received ", finishFlush)); ctx.setFlush(finishFlush); ctx.setInputBuffer(getInputFromSource(data)); return syncProcessBuffer(ctx, result); } template void ZlibUtil::brotliWithCallback( jsg::Lock& js, InputSource data, BrotliContext::Options options, CompressCallback cb) { // Capture only relevant errors so they can be passed to the callback auto res = js.tryCatch([&]() { return CompressCallbackArg(brotliSync(js, kj::mv(data), kj::mv(options))); }, [&](jsg::Value&& exception) { return CompressCallbackArg(jsg::JsValue(exception.getHandle(js))); }); // Ensure callback is invoked only once cb(js, kj::mv(res)); } template kj::Array ZlibUtil::zstdSync(jsg::Lock& js, InputSource data, ZstdContext::Options opts) { Context ctx(Context::Mode); auto chunkSize = opts.chunkSize.orDefault(ZLIB_PERFORMANT_CHUNK_SIZE); auto maxOutputLength = opts.maxOutputLength.orDefault(Z_MAX_CHUNK); JSG_REQUIRE(Z_MIN_CHUNK <= chunkSize && chunkSize <= Z_MAX_CHUNK, RangeError, kj::str("The value of \"options.chunkSize\" is out of range. It must be >= ", Z_MIN_CHUNK, " and <= ", Z_MAX_CHUNK, ". Received ", chunkSize)); JSG_REQUIRE(maxOutputLength >= 1 && maxOutputLength <= Z_MAX_CHUNK, RangeError, kj::str("The value of \"options.maxOutputLength\" is out of range. It must be >= 1 and <= ", Z_MAX_CHUNK, ". Received ", maxOutputLength)); GrowableBuffer result(ZLIB_PERFORMANT_CHUNK_SIZE, maxOutputLength); // Initialize the context if constexpr (Context::Mode == ZlibMode::ZSTD_ENCODE) { uint64_t pledgedSrcSize = opts.pledgedSrcSize.orDefault(ZSTD_CONTENTSIZE_UNKNOWN); KJ_IF_SOME(err, ctx.initialize(pledgedSrcSize)) { JSG_FAIL_REQUIRE(Error, err.message); } } else { KJ_IF_SOME(err, ctx.initialize()) { JSG_FAIL_REQUIRE(Error, err.message); } } // Set parameters KJ_IF_SOME(params, opts.params) { for (const auto& field: params.fields) { KJ_IF_SOME(err, ctx.setParams(field.name.parseAs(), field.value)) { JSG_FAIL_REQUIRE(Error, err.message); } } } auto flush = opts.flush.orDefault(ZSTD_e_continue); JSG_REQUIRE(ZSTD_e_continue <= flush && flush <= ZSTD_e_end, RangeError, kj::str("The value of \"options.flush\" is out of range. It must be >= ", ZSTD_e_continue, " and <= ", ZSTD_e_end, ". Received ", flush)); auto finishFlush = opts.finishFlush.orDefault(ZSTD_e_end); JSG_REQUIRE(ZSTD_e_continue <= finishFlush && finishFlush <= ZSTD_e_end, RangeError, kj::str("The value of \"options.finishFlush\" is out of range. It must be >= ", ZSTD_e_continue, " and <= ", ZSTD_e_end, ". Received ", finishFlush)); ctx.setFlush(finishFlush); ctx.setInputBuffer(getInputFromSource(data)); return syncProcessBuffer(ctx, result); } template void ZlibUtil::zstdWithCallback( jsg::Lock& js, InputSource data, ZstdContext::Options options, CompressCallback cb) { // Capture only relevant errors so they can be passed to the callback auto res = js.tryCatch([&]() { return CompressCallbackArg(zstdSync(js, kj::mv(data), kj::mv(options))); }, [&](jsg::Value&& exception) { return CompressCallbackArg(jsg::JsValue(exception.getHandle(js))); }); // Ensure callback is invoked only once cb(js, kj::mv(res)); } #ifndef CREATE_TEMPLATE #define CREATE_TEMPLATE(T) \ template class ZlibUtil::CompressionStream; \ template void ZlibUtil::CompressionStream::write(jsg::Lock & js, int flush, \ jsg::Optional> input, uint32_t inputOffset, uint32_t inputLength, \ kj::Array output, uint32_t outputOffset, uint32_t outputLength); \ template void ZlibUtil::CompressionStream::write(jsg::Lock & js, int flush, \ jsg::Optional> input, uint32_t inputOffset, uint32_t inputLength, \ kj::Array output, uint32_t outputOffset, uint32_t outputLength); CREATE_TEMPLATE(ZlibContext) CREATE_TEMPLATE(BrotliEncoderContext) CREATE_TEMPLATE(BrotliDecoderContext) CREATE_TEMPLATE(ZstdEncoderContext) CREATE_TEMPLATE(ZstdDecoderContext) template class ZlibUtil::BrotliCompressionStream; template class ZlibUtil::BrotliCompressionStream; template class ZlibUtil::ZstdCompressionStream; template class ZlibUtil::ZstdCompressionStream; template kj::Array ZlibUtil::brotliSync( jsg::Lock& js, InputSource data, BrotliContext::Options opts); template kj::Array ZlibUtil::brotliSync( jsg::Lock& js, InputSource data, BrotliContext::Options opts); template void ZlibUtil::brotliWithCallback( jsg::Lock& js, InputSource data, BrotliContext::Options options, CompressCallback cb); template void ZlibUtil::brotliWithCallback( jsg::Lock& js, InputSource data, BrotliContext::Options options, CompressCallback cb); template kj::Array ZlibUtil::zstdSync( jsg::Lock& js, InputSource data, ZstdContext::Options opts); template kj::Array ZlibUtil::zstdSync( jsg::Lock& js, InputSource data, ZstdContext::Options opts); template void ZlibUtil::zstdWithCallback( jsg::Lock& js, InputSource data, ZstdContext::Options options, CompressCallback cb); template void ZlibUtil::zstdWithCallback( jsg::Lock& js, InputSource data, ZstdContext::Options options, CompressCallback cb); #undef CREATE_TEMPLATE #endif } // namespace workerd::api::node