Skip to content
File

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

5.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 "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 
13namespace workerd::api {
14 
15namespace {
16template <typename T>
17jsg::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 
25jsg::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