Skip to content
File

Blob: src/workerd/api/streams/compression.c++

22.6 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "compression.h"
6 
7#include "nbytes.h"
8 
9#include <workerd/api/system-streams.h>
10#include <workerd/io/features.h>
11#include <workerd/util/autogate.h>
12#include <workerd/util/ring-buffer.h>
13#include <workerd/util/state-machine.h>
14 
15namespace workerd::api {
16CompressionAllocator::CompressionAllocator(
17 kj::Arc<const jsg::ExternalMemoryTarget>&& externalMemoryTarget)
18 : externalMemoryTarget(kj::mv(externalMemoryTarget)) {}
19 
20void CompressionAllocator::configure(z_stream* stream) {
21 stream->zalloc = AllocForZlib;
22 stream->zfree = FreeForZlib;
23 stream->opaque = this;
24}
25 
26void* CompressionAllocator::AllocForZlib(void* data, uInt items, uInt size) {
27 size_t real_size =
28 nbytes::MultiplyWithOverflowCheck(static_cast<size_t>(items), static_cast<size_t>(size));
29 return AllocForBrotli(data, real_size);
30}
31 
32void* CompressionAllocator::AllocForBrotli(void* opaque, size_t size) {
33 auto* allocator = static_cast<CompressionAllocator*>(opaque);
34 auto data = kj::heapArray<kj::byte>(size);
35 auto begin = data.begin();
36 
37 allocator->allocations.insert(begin,
38 {.data = kj::mv(data),
39 .memoryAdjustment = allocator->externalMemoryTarget->getAdjustment(size)});
40 return begin;
41}
42 
43void CompressionAllocator::FreeForZlib(void* opaque, void* pointer) {
44 if (KJ_UNLIKELY(pointer == nullptr)) return;
45 auto* allocator = static_cast<CompressionAllocator*>(opaque);
46 // No need to destroy memoryAdjustment here.
47 // Dropping the allocation from the hashmap will defer the adjustment
48 // until the isolate lock is held.
49 JSG_REQUIRE(allocator->allocations.erase(pointer), Error, "Zlib allocation should exist"_kj);
50}
51 
52namespace {
53 
54class Context {
55 public:
56 enum class Mode {
57 COMPRESS,
58 DECOMPRESS,
59 };
60 
61 enum class ContextFlags {
62 NONE,
63 STRICT,
64 };
65 
66 struct Result {
67 bool success = false;
68 kj::ArrayPtr<const byte> buffer;
69 };
70 
71 explicit Context(Mode mode,
72 kj::StringPtr format,
73 ContextFlags flags,
74 kj::Arc<const jsg::ExternalMemoryTarget>&& externalMemoryTarget)
75 : allocator(kj::mv(externalMemoryTarget)),
76 mode(mode),
77 strictCompression(flags)
78 
79 {
80 // Configure allocator before any stream operations.
81 allocator.configure(&ctx);
82 int result = Z_OK;
83 switch (mode) {
84 case Mode::COMPRESS:
85 result = deflateInit2(&ctx, Z_DEFAULT_COMPRESSION, Z_DEFLATED, getWindowBits(format),
86 8, // memLevel = 8 is the default
87 Z_DEFAULT_STRATEGY);
88 break;
89 case Mode::DECOMPRESS:
90 result = inflateInit2(&ctx, getWindowBits(format));
91 break;
92 default:
93 KJ_UNREACHABLE;
94 }
95 JSG_REQUIRE(result == Z_OK, Error, "Failed to initialize compression context."_kj);
96 }
97 
98 ~Context() noexcept(false) {
99 switch (mode) {
100 case Mode::COMPRESS:
101 deflateEnd(&ctx);
102 break;
103 case Mode::DECOMPRESS:
104 inflateEnd(&ctx);
105 break;
106 }
107 }
108 
109 KJ_DISALLOW_COPY_AND_MOVE(Context);
110 
111 void setInput(const void* in, size_t size) {
112 ctx.next_in = const_cast<byte*>(reinterpret_cast<const byte*>(in));
113 ctx.avail_in = size;
114 }
115 
116 Result pumpOnce(int flush) {
117 ctx.next_out = buffer;
118 ctx.avail_out = sizeof(buffer);
119 
120 int result = Z_OK;
121 
122 switch (mode) {
123 case Mode::COMPRESS:
124 result = deflate(&ctx, flush);
125 JSG_REQUIRE(result == Z_OK || result == Z_BUF_ERROR || result == Z_STREAM_END, TypeError,
126 "Compression failed.");
127 break;
128 case Mode::DECOMPRESS:
129 result = inflate(&ctx, flush);
130 JSG_REQUIRE(result == Z_OK || result == Z_BUF_ERROR || result == Z_STREAM_END, TypeError,
131 "Decompression failed.");
132 
133 if (strictCompression == ContextFlags::STRICT) {
134 // The spec requires that a TypeError is produced if there is trailing data after the end
135 // of the compression stream.
136 JSG_REQUIRE(!(result == Z_STREAM_END && ctx.avail_in > 0), TypeError,
137 "Trailing bytes after end of compressed data");
138 // Same applies to closing a stream before the complete decompressed data is available.
139 JSG_REQUIRE(
140 !(flush == Z_FINISH && result == Z_BUF_ERROR && ctx.avail_out == sizeof(buffer)),
141 TypeError, "Called close() on a decompression stream with incomplete data");
142 }
143 break;
144 default:
145 KJ_UNREACHABLE;
146 }
147 
148 return Result{
149 .success = result == Z_OK,
150 .buffer = kj::arrayPtr(buffer, sizeof(buffer) - ctx.avail_out),
151 };
152 }
153 
154 protected:
155 CompressionAllocator allocator;
156 
157 private:
158 static int getWindowBits(kj::StringPtr format) {
159 // We use a windowBits value of 15 combined with the magic value
160 // for the compression format type. For gzip, the magic value is
161 // 16, so the value returned is 15 + 16. For deflate, the magic
162 // value is 15. For raw deflate (i.e. deflate without a zlib header)
163 // the negative windowBits value is used, so -15. See the comments for
164 // deflateInit2() in zlib.h for details.
165 static constexpr auto GZIP = 16;
166 static constexpr auto DEFLATE = 15;
167 static constexpr auto DEFLATE_RAW = -15;
168 if (format == "gzip")
169 return DEFLATE + GZIP;
170 else if (format == "deflate")
171 return DEFLATE;
172 else if (format == "deflate-raw")
173 return DEFLATE_RAW;
174 KJ_UNREACHABLE;
175 }
176 
177 Mode mode;
178 z_stream ctx = {};
179 kj::byte buffer[16384];
180 
181 // For the eponymous compatibility flag
182 ContextFlags strictCompression;
183};
184 
185// Buffer class based on std::vector that erases data that has been read from it lazily to avoid
186// excessive copying when reading a larger amount of buffered data in small chunks. valid_size_ is
187// used to track the amount of data that has not been read back yet.
188class LazyBuffer {
189 public:
190 LazyBuffer(): valid_size_(0) {}
191 
192 // Return a chunk of data and mark it as invalid. The returned chunk remains valid until data is
193 // shifted, cleared or destructor is called. maybeShift() should be called after the returned data
194 // has been processed.
195 kj::ArrayPtr<byte> take(size_t read_size) {
196 KJ_ASSERT(read_size <= valid_size_);
197 kj::ArrayPtr<byte> chunk = kj::arrayPtr(&output[output.size() - valid_size_], read_size);
198 valid_size_ -= read_size;
199 return chunk;
200 }
201 
202 // Shift the output only if doing so results in reducing vector size by at least 1 KiB and 1/8 of
203 // its size to avoid copying for small reads.
204 void maybeShift() {
205 size_t unusedSpace = output.size() - valid_size_;
206 if (unusedSpace >= 1024 && unusedSpace >= (output.size() >> 3)) {
207 // Shifting buffer to erase data that has already been read. valid_size_ remains the same.
208 memmove(output.begin(), output.begin() + unusedSpace, valid_size_);
209 output.truncate(valid_size_);
210 }
211 }
212 
213 void write(kj::ArrayPtr<const byte> chunk) {
214 output.addAll(chunk);
215 valid_size_ += chunk.size();
216 }
217 
218 void clear() {
219 output.clear();
220 valid_size_ = 0;
221 }
222 
223 // For convenience, provide the size of the valid data that has not been read back yet. This may
224 // be smaller than the size of the internal vector, which is not relevant for the stream
225 // implementation.
226 size_t size() {
227 return valid_size_;
228 }
229 
230 // As with size(), the buffer is considered empty if there is no valid data remaining.
231 size_t empty() {
232 return valid_size_ == 0;
233 }
234 
235 private:
236 kj::Vector<kj::byte> output;
237 size_t valid_size_;
238};
239 
240// Because we have to use an autogate to switch things over to the new state manager, we need
241// to separate out a common base class for the compression stream internal state and separate
242// two separate impls that differ only in how they manage state. Once the autogate is removed,
243// we can delete the first impl class and merge everything back together.
244template <Context::Mode mode>
245class CompressionStreamBase: public kj::Refcounted,
246 public kj::AsyncInputStream,
247 public capnp::ExplicitEndOutputStream {
248 public:
249 explicit CompressionStreamBase(kj::String format,
250 Context::ContextFlags flags,
251 kj::Arc<const jsg::ExternalMemoryTarget>&& externalMemoryTarget)
252 : context(mode, format, flags, kj::mv(externalMemoryTarget)) {}
253 
254 // WritableStreamSink implementation ---------------------------------------------------
255 
256 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override final {
257 requireActive("Write after close");
258 context.setInput(buffer.begin(), buffer.size());
259 writeInternal(Z_NO_FLUSH);
260 co_return;
261 }
262 
263 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const kj::byte>> pieces) override final {
264 // We check state here so that we catch errors even if pieces is empty.
265 requireActive("Write after close");
266 for (auto piece: pieces) {
267 co_await write(piece);
268 }
269 co_return;
270 }
271 
272 kj::Promise<void> end() override final {
273 transitionToEnded();
274 writeInternal(Z_FINISH);
275 co_return;
276 }
277 
278 kj::Promise<void> whenWriteDisconnected() override final {
279 return kj::NEVER_DONE;
280 }
281 
282 void abortWrite(kj::Exception&& reason) override final {
283 cancelInternal(kj::mv(reason));
284 }
285 
286 // AsyncInputStream implementation -----------------------------------------------------
287 
288 kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override final {
289 KJ_ASSERT(minBytes <= maxBytes);
290 // Re-throw any stored exception
291 throwIfException();
292 // If stream has ended normally and no buffered data, return EOF
293 if (isInTerminalState() && output.empty()) {
294 co_return static_cast<size_t>(0);
295 }
296 // Active or terminal with data remaining
297 co_return co_await tryReadInternal(
298 kj::arrayPtr(reinterpret_cast<kj::byte*>(buffer), maxBytes), minBytes);
299 }
300 
301 protected:
302 virtual void requireActive(kj::StringPtr errorMessage) = 0;
303 virtual void transitionToEnded() = 0;
304 virtual void transitionToErrored(kj::Exception&& reason) = 0;
305 virtual void throwIfException() = 0;
306 virtual bool isInTerminalState() = 0;
307 
308 private:
309 struct PendingRead {
310 kj::ArrayPtr<kj::byte> buffer;
311 size_t minBytes = 1;
312 size_t filled = 0;
313 kj::Own<kj::PromiseFulfiller<size_t>> promise;
314 };
315 
316 void cancelInternal(kj::Exception reason) {
317 output.clear();
318 
319 while (!pendingReads.empty()) {
320 auto pending = kj::mv(pendingReads.front());
321 pendingReads.pop_front();
322 if (pending.promise->isWaiting()) {
323 pending.promise->reject(reason.clone());
324 }
325 }
326 
327 canceler.cancel(reason.clone());
328 transitionToErrored(kj::mv(reason));
329 }
330 
331 kj::Promise<size_t> tryReadInternal(kj::ArrayPtr<kj::byte> dest, size_t minBytes) {
332 const auto copyIntoBuffer = [this](kj::ArrayPtr<kj::byte> dest) {
333 auto maxBytesToCopy = kj::min(dest.size(), output.size());
334 dest.first(maxBytesToCopy).copyFrom(output.take(maxBytesToCopy));
335 output.maybeShift();
336 return maxBytesToCopy;
337 };
338 
339 // If the output currently contains >= minBytes, then we'll fulfill
340 // the read immediately, removing as many bytes as possible from the
341 // output queue.
342 // If we reached the end (terminal state), resolve the read immediately
343 // as well, since no new data is expected.
344 if (output.size() >= minBytes || isInTerminalState()) {
345 co_return copyIntoBuffer(dest);
346 }
347 
348 // Otherwise, create a pending read.
349 auto promise = kj::newPromiseAndFulfiller<size_t>();
350 auto pendingRead = PendingRead{
351 .buffer = dest,
352 .minBytes = minBytes,
353 .filled = 0,
354 .promise = kj::mv(promise.fulfiller),
355 };
356 
357 // If there are any bytes queued, copy as much as possible into the buffer.
358 if (output.size() > 0) {
359 pendingRead.filled = copyIntoBuffer(dest);
360 }
361 
362 pendingReads.push_back(kj::mv(pendingRead));
363 
364 co_return co_await canceler.wrap(kj::mv(promise.promise));
365 }
366 
367 void writeInternal(int flush) {
368 // TODO(later): This does not yet implement any backpressure. A caller can keep calling
369 // write without reading, which will continue to fill the internal buffer.
370 KJ_ASSERT(flush == Z_FINISH || !isInTerminalState());
371 Context::Result result;
372 
373 while (true) {
374 KJ_IF_SOME(exception, kj::runCatchingExceptions([this, flush, &result]() {
375 result = context.pumpOnce(flush);
376 })) {
377 cancelInternal(exception.clone());
378 kj::throwFatalException(kj::mv(exception));
379 }
380 
381 if (result.buffer.size() == 0) {
382 if (result.success) {
383 // No output produced but input data has been processed based on zlib return code, call
384 // pumpOnce again.
385 continue;
386 }
387 maybeFulfillRead();
388 return;
389 }
390 
391 // Output has been produced, copy it to result buffer and continue loop to call pumpOnce
392 // again.
393 output.write(result.buffer);
394 }
395 KJ_UNREACHABLE;
396 }
397 
398 // Fulfill as many pending reads as we can from the output buffer.
399 void maybeFulfillRead() {
400 // If there are pending reads and data to be read, we'll loop through
401 // the pending reads and fulfill them as much as possible.
402 while (!pendingReads.empty() && output.size() > 0) {
403 auto& pending = pendingReads.front();
404 
405 if (!pending.promise->isWaiting()) {
406 // The pending read was canceled!
407 // Importantly, the pending.buffer is no longer valid here so we definitely want to
408 // make sure we don't try to write anything to it!
409 
410 // If the pending read was already partially fulfilled, then we have a problem!
411 // We can't just cancel and continue because the partially read data will be lost
412 // so we need to report an error here and error the stream.
413 if (pending.filled > 0) {
414 auto ex = JSG_KJ_EXCEPTION(FAILED, Error, "A partially fulfilled read was canceled.");
415 cancelInternal(ex.clone());
416 kj::throwFatalException(kj::mv(ex));
417 }
418 
419 auto ex = JSG_KJ_EXCEPTION(FAILED, Error, "The pending read was canceled.");
420 cancelInternal(ex.clone());
421 kj::throwFatalException(kj::mv(ex));
422 }
423 
424 // The pending read is still viable so determine how much we can copy in.
425 auto amountToCopy = kj::min(pending.buffer.size() - pending.filled, output.size());
426 kj::ArrayPtr<byte> chunk = output.take(amountToCopy);
427 pending.buffer.slice(pending.filled, pending.filled + amountToCopy).copyFrom(chunk);
428 pending.filled += amountToCopy;
429 output.maybeShift();
430 
431 // If we've met the minimum bytes requirement for the pending read, fulfill
432 // the read promise.
433 if (pending.filled >= pending.minBytes) {
434 auto p = kj::mv(pending);
435 pendingReads.pop_front();
436 p.promise->fulfill(kj::mv(p.filled));
437 continue;
438 }
439 
440 // If we reached this point in the loop, remaining must be 0 so that we
441 // don't keep iterating through on the same pending read.
442 KJ_ASSERT(output.empty());
443 }
444 
445 if (isInTerminalState() && !pendingReads.empty()) {
446 // We are ended and we have pending reads. Because of the loop above,
447 // one of either pendingReads or output must be empty, so if we got this
448 // far, output.empty() must be true. Let's check.
449 KJ_ASSERT(output.empty());
450 // We need to flush any remaining reads.
451 while (!pendingReads.empty()) {
452 auto pending = kj::mv(pendingReads.front());
453 pendingReads.pop_front();
454 if (pending.promise->isWaiting()) {
455 // Fulfill the pending read promise only if it hasn't already been canceled.
456 pending.promise->fulfill(kj::mv(pending.filled));
457 }
458 }
459 }
460 }
461 
462 Context context;
463 
464 kj::Canceler canceler;
465 LazyBuffer output;
466 RingBuffer<PendingRead, 8> pendingReads;
467};
468 
469template <Context::Mode mode>
470class CompressionStreamImpl final: public CompressionStreamBase<mode> {
471 public:
472 explicit CompressionStreamImpl(kj::String format,
473 Context::ContextFlags flags,
474 kj::Arc<const jsg::ExternalMemoryTarget>&& externalMemoryTarget)
475 : CompressionStreamBase<mode>(kj::mv(format), flags, kj::mv(externalMemoryTarget)),
476 state(decltype(state)::template create<Open>()) {}
477 
478 protected:
479 void requireActive(kj::StringPtr errorMessage) override {
480 KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) {
481 kj::throwFatalException(exception.clone());
482 }
483 // isActive() returns true only if in Open state (the ActiveState)
484 JSG_REQUIRE(state.isActive(), Error, errorMessage);
485 }
486 
487 void transitionToEnded() override {
488 // If already in a terminal state (Ended or Exception), this is a no-op.
489 // This matches the V1 behavior where calling end() multiple times was allowed.
490 if (state.isTerminal()) return;
491 auto result = state.template transitionFromTo<Open, Ended>();
492 KJ_REQUIRE(result != kj::none, "Stream already ended or errored");
493 }
494 
495 void transitionToErrored(kj::Exception&& reason) override {
496 // Use forceTransitionTo because cancelInternal may be called when already
497 // in an error state (e.g., from writeInternal error handling).
498 state.template forceTransitionTo<kj::Exception>(kj::mv(reason));
499 }
500 
501 void throwIfException() override {
502 KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) {
503 kj::throwFatalException(exception.clone());
504 }
505 }
506 
507 virtual bool isInTerminalState() override {
508 return state.isTerminal();
509 }
510 
511 private:
512 struct Ended {
513 static constexpr kj::StringPtr NAME KJ_UNUSED = "ended"_kj;
514 };
515 struct Open {
516 static constexpr kj::StringPtr NAME KJ_UNUSED = "open"_kj;
517 };
518 
519 // State machine for tracking compression stream lifecycle:
520 // Open -> Ended (normal close via end())
521 // Open -> kj::Exception (error via abortWrite())
522 // Ended is terminal, kj::Exception is implicitly terminal via ErrorState.
523 StateMachine<TerminalStates<Ended>,
524 ErrorState<kj::Exception>,
525 ActiveState<Open>,
526 Open,
527 Ended,
528 kj::Exception>
529 state;
530};
531 
532// Adapter to bridge CompressionStreamImpl (which implements AsyncInputStream and
533// ExplicitEndOutputStream) to the ReadableStreamSource/WritableStreamSink interfaces.
534// TODO(soon): This class is intended to be replaced by the new ReadableSource/WritableSink
535// interfaces once fully implemented. We will need an adapter that knows how to handle both
536// sides of the stream once fully implemented. The current implementation in system-streams.c++
537// implements separate adapters for each side that are not aware of each other, making it
538// unsuitable for this specific case.
539template <Context::Mode mode>
540class CompressionStreamAdapter final: public kj::Refcounted,
541 public ReadableStreamSource,
542 public WritableStreamSink {
543 public:
544 explicit CompressionStreamAdapter(kj::Rc<CompressionStreamBase<mode>> impl)
545 : impl(kj::mv(impl)),
546 ioContext(IoContext::current()) {}
547 
548 // ReadableStreamSource implementation
549 kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override {
550 return impl->tryRead(buffer, minBytes, maxBytes).attach(ioContext.registerPendingEvent());
551 }
552 
553 void cancel(kj::Exception reason) override {
554 // AsyncInputStream doesn't have cancel, but we can abort the write side
555 impl->abortWrite(kj::mv(reason));
556 }
557 
558 // WritableStreamSink implementation
559 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
560 return impl->write(buffer).attach(ioContext.registerPendingEvent());
561 }
562 
563 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
564 return impl->write(pieces).attach(ioContext.registerPendingEvent());
565 }
566 
567 kj::Promise<void> end() override {
568 return impl->end().attach(ioContext.registerPendingEvent());
569 }
570 
571 void abort(kj::Exception reason) override {
572 impl->abortWrite(kj::mv(reason));
573 }
574 
575 private:
576 kj::Rc<CompressionStreamBase<mode>> impl;
577 IoContext& ioContext;
578};
579 
580kj::Rc<CompressionStreamBase<Context::Mode::COMPRESS>> createCompressionStreamImpl(
581 kj::String format,
582 Context::ContextFlags flags,
583 kj::Arc<const jsg::ExternalMemoryTarget>&& externalMemoryTarget) {
584 return kj::rc<CompressionStreamImpl<Context::Mode::COMPRESS>>(
585 kj::mv(format), flags, kj::mv(externalMemoryTarget));
586}
587 
588kj::Rc<CompressionStreamBase<Context::Mode::DECOMPRESS>> createDecompressionStreamImpl(
589 kj::String format,
590 Context::ContextFlags flags,
591 kj::Arc<const jsg::ExternalMemoryTarget>&& externalMemoryTarget) {
592 return kj::rc<CompressionStreamImpl<Context::Mode::DECOMPRESS>>(
593 kj::mv(format), flags, kj::mv(externalMemoryTarget));
594}
595 
596} // namespace
597 
598jsg::Ref<CompressionStream> CompressionStream::constructor(jsg::Lock& js, kj::String format) {
599 JSG_REQUIRE(format == "deflate" || format == "gzip" || format == "deflate-raw", TypeError,
600 "The compression format must be either 'deflate', 'deflate-raw' or 'gzip'.");
601 
602 // TODO(cleanup): Once the autogate is removed, we can delete CompressionStreamImpl
603 kj::Rc<CompressionStreamBase<Context::Mode::COMPRESS>> impl = createCompressionStreamImpl(
604 kj::mv(format), Context::ContextFlags::NONE, js.getExternalMemoryTarget());
605 
606 auto& ioContext = IoContext::current();
607 
608 // Create a single adapter that implements both readable and writable sides
609 auto adapter = kj::refcounted<CompressionStreamAdapter<Context::Mode::COMPRESS>>(kj::mv(impl));
610 auto readableSide = kj::addRef(*adapter);
611 auto writableSide = kj::mv(adapter);
612 
613 return js.alloc<CompressionStream>(js.alloc<ReadableStream>(ioContext, kj::mv(readableSide)),
614 js.alloc<WritableStream>(ioContext, kj::mv(writableSide),
615 ioContext.getMetrics().tryCreateWritableByteStreamObserver()));
616}
617 
618jsg::Ref<DecompressionStream> DecompressionStream::constructor(jsg::Lock& js, kj::String format) {
619 JSG_REQUIRE(format == "deflate" || format == "gzip" || format == "deflate-raw", TypeError,
620 "The compression format must be either 'deflate', 'deflate-raw' or 'gzip'.");
621 
622 kj::Rc<CompressionStreamBase<Context::Mode::DECOMPRESS>> impl =
623 createDecompressionStreamImpl(kj::mv(format),
624 FeatureFlags::get(js).getStrictCompression() ? Context::ContextFlags::STRICT
625 : Context::ContextFlags::NONE,
626 js.getExternalMemoryTarget());
627 
628 auto& ioContext = IoContext::current();
629 
630 // Create a single adapter that implements both readable and writable sides
631 auto adapter = kj::refcounted<CompressionStreamAdapter<Context::Mode::DECOMPRESS>>(kj::mv(impl));
632 auto readableSide = kj::addRef(*adapter);
633 auto writableSide = kj::mv(adapter);
634 
635 return js.alloc<DecompressionStream>(js.alloc<ReadableStream>(ioContext, kj::mv(readableSide)),
636 js.alloc<WritableStream>(ioContext, kj::mv(writableSide),
637 ioContext.getMetrics().tryCreateWritableByteStreamObserver()));
638}
639 
640} // namespace workerd::api