File
Blob: src/workerd/api/node/zlib-util.c++
| 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 |
| 14 | namespace workerd::api::node { |
| 15 | |
| 16 | kj::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 | |
| 32 | uint32_t ZlibUtil::crc32Sync(InputSource data, uint32_t value) { |
| 33 | auto dataPtr = getInputFromSource(data); |
| 34 | return crc32(value, dataPtr.begin(), dataPtr.size()); |
| 35 | } |
| 36 | |
| 37 | namespace { |
| 38 | class 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 | |
| 121 | void 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 | |
| 171 | kj::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 | |
| 197 | kj::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 | |
| 223 | bool 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 | |
| 254 | kj::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 | |
| 282 | void 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 | |
| 381 | kj::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 | |
| 404 | ZlibContext::~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 | |
| 430 | void 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 | |
| 437 | void 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 | |
| 443 | void ZlibContext::setOutputBuffer(kj::ArrayPtr<kj::byte> output) { |
| 444 | stream.next_out = output.begin(); |
| 445 | stream.avail_out = output.size(); |
| 446 | } |
| 447 | |
| 448 | template <typename CompressionContext> |
| 449 | jsg::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 | |
| 454 | template <typename CompressionContext> |
| 455 | ZlibUtil::CompressionStream<CompressionContext>::~CompressionStream() { |
| 456 | JSG_ASSERT(!writing, Error, "Writing to compression stream"_kj); |
| 457 | close(); |
| 458 | } |
| 459 | |
| 460 | template <typename CompressionContext> |
| 461 | void 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 | |
| 473 | template <typename CompressionContext> |
| 474 | template <bool async> |
| 475 | void 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 | |
| 532 | template <typename CompressionContext> |
| 533 | void 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 | |
| 543 | template <typename CompressionContext> |
| 544 | bool 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 | |
| 552 | template <typename CompressionContext> |
| 553 | void 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 | |
| 560 | template <typename CompressionContext> |
| 561 | void 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 | |
| 570 | template <typename CompressionContext> |
| 571 | template <bool async> |
| 572 | void 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 | |
| 606 | template <typename CompressionContext> |
| 607 | void ZlibUtil::CompressionStream<CompressionContext>::reset(jsg::Lock& js) { |
| 608 | KJ_IF_SOME(error, context()->resetStream()) { |
| 609 | emitError(js, kj::mv(error)); |
| 610 | } |
| 611 | } |
| 612 | |
| 613 | jsg::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 | |
| 618 | void 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 | |
| 631 | void 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 | |
| 638 | void 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 | |
| 645 | void BrotliContext::setInputBuffer(kj::ArrayPtr<const kj::byte> input) { |
| 646 | nextIn = input.begin(); |
| 647 | availIn = input.size(); |
| 648 | } |
| 649 | |
| 650 | void BrotliContext::setOutputBuffer(kj::ArrayPtr<kj::byte> output) { |
| 651 | nextOut = output.begin(); |
| 652 | availOut = output.size(); |
| 653 | } |
| 654 | |
| 655 | uint BrotliContext::getAvailOut() const { |
| 656 | return availOut; |
| 657 | } |
| 658 | |
| 659 | void BrotliContext::setFlush(int _flush) { |
| 660 | flush = static_cast<BrotliEncoderOperation>(_flush); |
| 661 | } |
| 662 | |
| 663 | void BrotliContext::getAfterWriteResult(uint32_t* _availIn, uint32_t* _availOut) const { |
| 664 | *_availIn = availIn; |
| 665 | *_availOut = availOut; |
| 666 | } |
| 667 | |
| 668 | BrotliEncoderContext::BrotliEncoderContext(ZlibMode _mode): BrotliContext(_mode) { |
| 669 | auto instance = BrotliEncoderCreateInstance(alloc_brotli, free_brotli, alloc_opaque_brotli); |
| 670 | state = kj::disposeWith<BrotliEncoderDestroyInstance>(instance); |
| 671 | } |
| 672 | |
| 673 | void 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 | |
| 685 | kj::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 | |
| 702 | kj::Maybe<CompressionError> BrotliEncoderContext::resetStream() { |
| 703 | return initialize(alloc_brotli, free_brotli, alloc_opaque_brotli); |
| 704 | } |
| 705 | |
| 706 | kj::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 | |
| 714 | kj::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 | |
| 722 | bool BrotliEncoderContext::isStreamEnd() const { |
| 723 | return streamEnd; |
| 724 | } |
| 725 | |
| 726 | BrotliDecoderContext::BrotliDecoderContext(ZlibMode _mode): BrotliContext(_mode) { |
| 727 | auto instance = BrotliDecoderCreateInstance(alloc_brotli, free_brotli, alloc_opaque_brotli); |
| 728 | state = kj::disposeWith<BrotliDecoderDestroyInstance>(instance); |
| 729 | } |
| 730 | |
| 731 | kj::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 | |
| 748 | void 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 | |
| 762 | kj::Maybe<CompressionError> BrotliDecoderContext::resetStream() { |
| 763 | return initialize(alloc_brotli, free_brotli, alloc_opaque_brotli); |
| 764 | } |
| 765 | |
| 766 | kj::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 | |
| 774 | kj::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 | |
| 787 | bool BrotliDecoderContext::isStreamEnd() const { |
| 788 | return lastResult == BROTLI_DECODER_RESULT_SUCCESS; |
| 789 | } |
| 790 | |
| 791 | // ======================================================================================= |
| 792 | // Zstd Implementation |
| 793 | |
| 794 | void ZstdContext::setBuffers(kj::ArrayPtr<kj::byte> input, kj::ArrayPtr<kj::byte> output) { |
| 795 | setInputBuffer(input); |
| 796 | setOutputBuffer(output); |
| 797 | } |
| 798 | |
| 799 | void ZstdContext::setInputBuffer(kj::ArrayPtr<const kj::byte> input) { |
| 800 | input_.src = input.begin(); |
| 801 | input_.size = input.size(); |
| 802 | input_.pos = 0; |
| 803 | } |
| 804 | |
| 805 | void ZstdContext::setOutputBuffer(kj::ArrayPtr<kj::byte> output) { |
| 806 | output_.dst = output.begin(); |
| 807 | output_.size = output.size(); |
| 808 | output_.pos = 0; |
| 809 | } |
| 810 | |
| 811 | void 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 | |
| 817 | kj::uint ZstdContext::getAvailOut() const { |
| 818 | return output_.size - output_.pos; |
| 819 | } |
| 820 | |
| 821 | void ZstdContext::getAfterWriteResult(uint32_t* availIn, uint32_t* availOut) const { |
| 822 | *availIn = input_.size - input_.pos; |
| 823 | *availOut = output_.size - output_.pos; |
| 824 | } |
| 825 | |
| 826 | namespace { |
| 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. |
| 829 | kj::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). |
| 839 | void zstdFreeCCtx(ZSTD_CCtx* cctx) { |
| 840 | ZSTD_freeCCtx(cctx); |
| 841 | } |
| 842 | void zstdFreeDCtx(ZSTD_DCtx* dctx) { |
| 843 | ZSTD_freeDCtx(dctx); |
| 844 | } |
| 845 | } // namespace |
| 846 | |
| 847 | ZstdEncoderContext::ZstdEncoderContext(ZlibMode _mode) |
| 848 | : ZstdContext(_mode), |
| 849 | cctx_(kj::disposeWith<zstdFreeCCtx>(ZSTD_createCCtx())) {} |
| 850 | |
| 851 | kj::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 | |
| 867 | void 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 | |
| 878 | kj::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 | |
| 888 | kj::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 | |
| 899 | kj::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 | |
| 913 | bool 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 | |
| 918 | ZstdDecoderContext::ZstdDecoderContext(ZlibMode _mode) |
| 919 | : ZstdContext(_mode), |
| 920 | dctx_(kj::disposeWith<zstdFreeDCtx>(ZSTD_createDCtx())) {} |
| 921 | |
| 922 | kj::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 | |
| 933 | void 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 | |
| 948 | kj::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 | |
| 959 | kj::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 | |
| 969 | kj::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 | |
| 985 | bool 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 | |
| 990 | template <typename CompressionContext> |
| 991 | jsg::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 | |
| 996 | template <typename CompressionContext> |
| 997 | bool 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 | |
| 1033 | template <typename CompressionContext> |
| 1034 | jsg::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 | |
| 1040 | template <typename CompressionContext> |
| 1041 | bool 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 | |
| 1069 | namespace { |
| 1070 | template <typename Context> |
| 1071 | static 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 | |
| 1097 | kj::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 | |
| 1136 | void 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 | |
| 1152 | template <typename Context> |
| 1153 | kj::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 | |
| 1206 | template <typename Context> |
| 1207 | void 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 | |
| 1220 | template <typename Context> |
| 1221 | kj::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 | |
| 1271 | template <typename Context> |
| 1272 | void 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 | |
| 1295 | CREATE_TEMPLATE(ZlibContext) |
| 1296 | CREATE_TEMPLATE(BrotliEncoderContext) |
| 1297 | CREATE_TEMPLATE(BrotliDecoderContext) |
| 1298 | CREATE_TEMPLATE(ZstdEncoderContext) |
| 1299 | CREATE_TEMPLATE(ZstdDecoderContext) |
| 1300 | |
| 1301 | template class ZlibUtil::BrotliCompressionStream<BrotliEncoderContext>; |
| 1302 | template class ZlibUtil::BrotliCompressionStream<BrotliDecoderContext>; |
| 1303 | |
| 1304 | template class ZlibUtil::ZstdCompressionStream<ZstdEncoderContext>; |
| 1305 | template class ZlibUtil::ZstdCompressionStream<ZstdDecoderContext>; |
| 1306 | |
| 1307 | template kj::Array<kj::byte> ZlibUtil::brotliSync<BrotliEncoderContext>( |
| 1308 | jsg::Lock& js, InputSource data, BrotliContext::Options opts); |
| 1309 | template kj::Array<kj::byte> ZlibUtil::brotliSync<BrotliDecoderContext>( |
| 1310 | jsg::Lock& js, InputSource data, BrotliContext::Options opts); |
| 1311 | template void ZlibUtil::brotliWithCallback<BrotliEncoderContext>( |
| 1312 | jsg::Lock& js, InputSource data, BrotliContext::Options options, CompressCallback cb); |
| 1313 | template void ZlibUtil::brotliWithCallback<BrotliDecoderContext>( |
| 1314 | jsg::Lock& js, InputSource data, BrotliContext::Options options, CompressCallback cb); |
| 1315 | |
| 1316 | template kj::Array<kj::byte> ZlibUtil::zstdSync<ZstdEncoderContext>( |
| 1317 | jsg::Lock& js, InputSource data, ZstdContext::Options opts); |
| 1318 | template kj::Array<kj::byte> ZlibUtil::zstdSync<ZstdDecoderContext>( |
| 1319 | jsg::Lock& js, InputSource data, ZstdContext::Options opts); |
| 1320 | template void ZlibUtil::zstdWithCallback<ZstdEncoderContext>( |
| 1321 | jsg::Lock& js, InputSource data, ZstdContext::Options options, CompressCallback cb); |
| 1322 | template 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 |