File
Blob: src/workerd/api/streams/transform.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 | |
| 5 | #include "transform.h" |
| 6 | |
| 7 | #include "identity-transform-stream.h" |
| 8 | #include "standard.h" |
| 9 | |
| 10 | #include <workerd/io/features.h> |
| 11 | #include <workerd/jsg/jsg.h> |
| 12 | |
| 13 | namespace workerd::api { |
| 14 | |
| 15 | namespace { |
| 16 | template <typename T> |
| 17 | jsg::Function<T> maybeAddFunctor(auto t) { |
| 18 | KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { |
| 19 | return jsg::Function<T>(ioContext.addFunctor(kj::mv(t))); |
| 20 | } |
| 21 | return jsg::Function<T>(kj::mv(t)); |
| 22 | } |
| 23 | } // namespace |
| 24 | |
| 25 | jsg::Ref<TransformStream> TransformStream::constructor(jsg::Lock& js, |
| 26 | jsg::Optional<Transformer> maybeTransformer, |
| 27 | jsg::Optional<StreamQueuingStrategy> maybeWritableStrategy, |
| 28 | jsg::Optional<StreamQueuingStrategy> maybeReadableStrategy) { |
| 29 | |
| 30 | if (FeatureFlags::get(js).getTransformStreamJavaScriptControllers()) { |
| 31 | // The standard implementation. Here the TransformStream is backed by readable |
| 32 | // and writable streams using the JavaScript-backed controllers. Data that is |
| 33 | // written to the writable side passes through the transform function that is |
| 34 | // given in maybeTransformer. If no transform function is given, then any value |
| 35 | // written is passed through unchanged. |
| 36 | // |
| 37 | // Per the standard specification, any JavaScript value can be written to and |
| 38 | // read from the transform stream, and the readable side does *not* support BYOB |
| 39 | // reads. |
| 40 | // |
| 41 | // Persistent references to the TransformStreamDefaultController are held by both |
| 42 | // the readable and writable sides. The actual TransformStream object can be dropped |
| 43 | // and allowed to be garbage collected. |
| 44 | |
| 45 | auto controller = js.alloc<TransformStreamDefaultController>(js); |
| 46 | auto transformer = kj::mv(maybeTransformer).orDefault({}); |
| 47 | |
| 48 | // By default, let's signal backpressure on the readable side by setting the highWaterMark |
| 49 | // to zero if a strategy is not given. This effectively means that writes/reads will be |
| 50 | // one to one as long as the writer is respecting backpressure signals. If buffering |
| 51 | // occurs, it will happen in the writable side of the transform stream. |
| 52 | auto readableStrategy = kj::mv(maybeReadableStrategy) |
| 53 | .orDefault(StreamQueuingStrategy{ |
| 54 | .highWaterMark = 0, |
| 55 | }); |
| 56 | |
| 57 | auto readable = ReadableStream::constructor(js, |
| 58 | UnderlyingSource{ |
| 59 | .type = kj::none, |
| 60 | .autoAllocateChunkSize = kj::none, |
| 61 | .start = maybeAddFunctor<UnderlyingSource::StartAlgorithm>( |
| 62 | JSG_VISITABLE_LAMBDA((controller = controller.addRef()), (controller), |
| 63 | (jsg::Lock & js, auto c) mutable { return controller->getStartPromise(js); })), |
| 64 | .pull = maybeAddFunctor<UnderlyingSource::PullAlgorithm>( |
| 65 | JSG_VISITABLE_LAMBDA((controller = controller.addRef()), (controller), |
| 66 | (jsg::Lock & js, auto c) mutable { return controller->pull(js); })), |
| 67 | .cancel = maybeAddFunctor<UnderlyingSource::CancelAlgorithm>( JSG_VISITABLE_LAMBDA( |
| 68 | (controller = controller.addRef()), (controller), |
| 69 | (jsg::Lock & js, auto reason) mutable { return controller->cancel(js, reason); })), |
| 70 | .expectedLength = transformer.expectedLength.map( |
| 71 | [](uint64_t expectedLength) { return expectedLength; }), |
| 72 | }, |
| 73 | kj::mv(readableStrategy)); |
| 74 | |
| 75 | auto writable = WritableStream::constructor(js, |
| 76 | UnderlyingSink{ |
| 77 | .type = kj::none, |
| 78 | .start = maybeAddFunctor<UnderlyingSink::StartAlgorithm>( |
| 79 | JSG_VISITABLE_LAMBDA((controller = controller.addRef()), (controller), |
| 80 | (jsg::Lock & js, auto c) mutable { return controller->getStartPromise(js); })), |
| 81 | .write = maybeAddFunctor<UnderlyingSink::WriteAlgorithm>( |
| 82 | JSG_VISITABLE_LAMBDA((controller = controller.addRef()), (controller), |
| 83 | (jsg::Lock & js, auto chunk, auto c) mutable { |
| 84 | return controller->write(js, chunk); |
| 85 | })), |
| 86 | .abort = maybeAddFunctor<UnderlyingSink::AbortAlgorithm>( |
| 87 | JSG_VISITABLE_LAMBDA((controller = controller.addRef()), (controller), |
| 88 | (jsg::Lock & js, auto reason) mutable { return controller->abort(js, reason); })), |
| 89 | .close = maybeAddFunctor<UnderlyingSink::CloseAlgorithm>( |
| 90 | JSG_VISITABLE_LAMBDA((controller = controller.addRef()), (controller), |
| 91 | (jsg::Lock & js) mutable { return controller->close(js); })), |
| 92 | }, |
| 93 | kj::mv(maybeWritableStrategy)); |
| 94 | |
| 95 | // The controller will store c++ references to both the readable and writable |
| 96 | // streams underlying controllers. |
| 97 | controller->init(js, readable, writable, kj::mv(transformer)); |
| 98 | |
| 99 | return js.alloc<TransformStream>(kj::mv(readable), kj::mv(writable)); |
| 100 | } |
| 101 | |
| 102 | // The old implementation just defers to IdentityTransformStream. If any of the arguments |
| 103 | // are specified we throw because it's most likely that they want the standard implementation |
| 104 | // but the compatibility flag is not set. |
| 105 | if (maybeTransformer != kj::none || maybeWritableStrategy != kj::none || |
| 106 | maybeReadableStrategy != kj::none) { |
| 107 | IoContext::current().logWarningOnce( |
| 108 | "To use the new TransformStream() constructor with a " |
| 109 | "custom transformer, enable the transformstream_enable_standard_constructor compatibility flag. " |
| 110 | "Refer to the docs for more information: https://developers.cloudflare.com/workers/platform/compatibility-dates/#compatibility-flags"); |
| 111 | } |
| 112 | |
| 113 | return IdentityTransformStream::constructor(js); |
| 114 | } |
| 115 | |
| 116 | } // namespace workerd::api |