File
Blob: src/workerd/api/streams/common.h
| 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 | #pragma once |
| 6 | |
| 7 | #include "../basics.h" |
| 8 | |
| 9 | #include <workerd/io/io-context.h> |
| 10 | #include <workerd/io/worker-interface.capnp.h> |
| 11 | #include <workerd/jsg/jsg.h> |
| 12 | |
| 13 | #if _MSC_VER |
| 14 | using ssize_t = long long; |
| 15 | #endif |
| 16 | |
| 17 | namespace workerd::api { |
| 18 | |
| 19 | class ReadableStream; |
| 20 | class ReadableStreamController; |
| 21 | class ReadableStreamSource; |
| 22 | class ReadableStreamDefaultController; |
| 23 | class ReadableByteStreamController; |
| 24 | |
| 25 | class WritableStream; |
| 26 | class WritableStreamController; |
| 27 | class WritableStreamSink; |
| 28 | class WritableStreamDefaultController; |
| 29 | |
| 30 | class TransformStreamDefaultController; |
| 31 | |
| 32 | using rpc::StreamEncoding; |
| 33 | |
| 34 | enum class ReadAllTextOption : uint8_t { |
| 35 | NONE = 0, |
| 36 | NULL_TERMINATE = 1 << 0, |
| 37 | STRIP_BOM = 1 << 1, |
| 38 | }; |
| 39 | |
| 40 | inline ReadAllTextOption operator|(ReadAllTextOption a, ReadAllTextOption b) { |
| 41 | return static_cast<ReadAllTextOption>(static_cast<uint8_t>(a) | static_cast<uint8_t>(b)); |
| 42 | } |
| 43 | |
| 44 | inline ReadAllTextOption& operator|=(ReadAllTextOption& a, ReadAllTextOption b) { |
| 45 | return a = a | b; |
| 46 | } |
| 47 | |
| 48 | inline bool operator&(ReadAllTextOption a, ReadAllTextOption b) { |
| 49 | return (static_cast<uint8_t>(a) & static_cast<uint8_t>(b)) != 0; |
| 50 | } |
| 51 | |
| 52 | static constexpr kj::byte UTF8_BOM[] = {0xEF, 0xBB, 0xBF}; |
| 53 | static constexpr size_t UTF8_BOM_SIZE = sizeof(UTF8_BOM); |
| 54 | |
| 55 | inline bool hasUtf8Bom(kj::ArrayPtr<const kj::byte> data) { |
| 56 | return data.size() >= UTF8_BOM_SIZE && memcmp(data.begin(), UTF8_BOM, UTF8_BOM_SIZE) == 0; |
| 57 | } |
| 58 | |
| 59 | struct ReadResult { |
| 60 | jsg::Optional<jsg::Value> value; |
| 61 | bool done; |
| 62 | |
| 63 | JSG_STRUCT(value, done); |
| 64 | JSG_STRUCT_TS_OVERRIDE(type ReadableStreamReadResult<R = any> = |
| 65 | | { done: false, value: R; } |
| 66 | | { done: true; value?: undefined; } |
| 67 | ); |
| 68 | |
| 69 | void visitForGc(jsg::GcVisitor& visitor) { |
| 70 | visitor.visit(value); |
| 71 | } |
| 72 | }; |
| 73 | |
| 74 | // Result type for draining read operations. Always returns bytes, even for value streams. |
| 75 | // Used by DrainingReader for optimized pipe-to operations with vectored writes. |
| 76 | // This is a C++ only type - not exposed to JavaScript. |
| 77 | struct DrainingReadResult { |
| 78 | kj::Array<kj::Array<kj::byte>> chunks; // Multiple byte arrays for vectored writes |
| 79 | bool done = false; // True if stream is closed/closing |
| 80 | }; |
| 81 | |
| 82 | struct StreamQueuingStrategy { |
| 83 | using SizeAlgorithm = uint64_t(v8::Local<v8::Value>); |
| 84 | |
| 85 | jsg::Optional<uint64_t> highWaterMark; |
| 86 | jsg::Optional<jsg::Function<SizeAlgorithm>> size; |
| 87 | |
| 88 | JSG_STRUCT(highWaterMark, size); |
| 89 | JSG_STRUCT_TS_OVERRIDE(QueuingStrategy<T = any> { |
| 90 | size?: (chunk: T) => number | bigint; |
| 91 | }); |
| 92 | }; |
| 93 | |
| 94 | struct UnderlyingSource { |
| 95 | using Controller = |
| 96 | kj::OneOf<jsg::Ref<ReadableStreamDefaultController>, jsg::Ref<ReadableByteStreamController>>; |
| 97 | using StartAlgorithm = jsg::Promise<void>(Controller); |
| 98 | using PullAlgorithm = jsg::Promise<void>(Controller); |
| 99 | using CancelAlgorithm = jsg::Promise<void>(v8::Local<v8::Value> reason); |
| 100 | |
| 101 | // The autoAllocateChunkSize mechanism allows byte streams to operate as if a BYOB |
| 102 | // reader is being used even if it is just a default reader. Support is optional |
| 103 | // per the streams spec but our implementation will always enable it. Specifically, |
| 104 | // if user code does not provide an explicit autoAllocateChunkSize, we'll assume |
| 105 | // this default. |
| 106 | static constexpr int DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE = 4096; |
| 107 | |
| 108 | // We want to increase the default auto allocate chunk size but we need to do |
| 109 | // so carefully to avoid introducing memory regressions and causing workers to |
| 110 | // hit OOM errors. We'll use an autogate to roll out the new default. |
| 111 | static constexpr int DEFAULT_AUTO_ALLOCATE_CHUNK_SIZE_2 = 16 * 1024; |
| 112 | |
| 113 | // Per the spec, the type property for the UnderlyingSource should be either |
| 114 | // undefined, the empty string, or "bytes". When undefined, the empty string is |
| 115 | // used as the default. When type is the empty string, the stream is considered |
| 116 | // to be value-oriented rather than byte-oriented. |
| 117 | jsg::Optional<kj::String> type; |
| 118 | |
| 119 | // Used only when type is equal to "bytes", the autoAllocateChunkSize defines |
| 120 | // the size of automatically allocated buffer that is created when a default |
| 121 | // mode read is performed on a byte-oriented ReadableStream that supports |
| 122 | // BYOB reads. The stream standard makes this optional to support and defines |
| 123 | // no default value. We've chosen to use a default value of 4096. If given, |
| 124 | // the value must be greater than zero. |
| 125 | jsg::Optional<int> autoAllocateChunkSize; |
| 126 | |
| 127 | jsg::Optional<jsg::Function<StartAlgorithm>> start; |
| 128 | jsg::Optional<jsg::Function<PullAlgorithm>> pull; |
| 129 | jsg::Optional<jsg::Function<CancelAlgorithm>> cancel; |
| 130 | |
| 131 | // The expectedLength is a non-standard extension used to support specifying the |
| 132 | // content-length when using a ReadableStream as the body of a request or response. |
| 133 | jsg::Optional<uint64_t> expectedLength; |
| 134 | |
| 135 | JSG_STRUCT(type, autoAllocateChunkSize, start, pull, cancel, expectedLength); |
| 136 | JSG_STRUCT_TS_DEFINE(interface UnderlyingByteSource { |
| 137 | type: "bytes"; |
| 138 | autoAllocateChunkSize?: number; |
| 139 | start?: (controller: ReadableByteStreamController) => void | Promise<void>; |
| 140 | pull?: (controller: ReadableByteStreamController) => void | Promise<void>; |
| 141 | cancel?: (reason: any) => void | Promise<void>; |
| 142 | }); |
| 143 | JSG_STRUCT_TS_OVERRIDE(<R = any> { |
| 144 | type?: "" | undefined; |
| 145 | autoAllocateChunkSize: never; |
| 146 | start?: (controller: ReadableStreamDefaultController<R>) => void | Promise<void>; |
| 147 | pull?: (controller: ReadableStreamDefaultController<R>) => void | Promise<void>; |
| 148 | cancel?: (reason: any) => void | Promise<void>; |
| 149 | }); |
| 150 | }; |
| 151 | |
| 152 | struct UnderlyingSink { |
| 153 | using Controller = jsg::Ref<WritableStreamDefaultController>; |
| 154 | using StartAlgorithm = jsg::Promise<void>(Controller); |
| 155 | using WriteAlgorithm = jsg::Promise<void>(v8::Local<v8::Value>, Controller); |
| 156 | using AbortAlgorithm = jsg::Promise<void>(v8::Local<v8::Value> reason); |
| 157 | using CloseAlgorithm = jsg::Promise<void>(); |
| 158 | |
| 159 | // Per the spec, the type property for the UnderlyingSink should always be either |
| 160 | // undefined or the empty string. Any other value will trigger a TypeError. |
| 161 | jsg::Optional<kj::String> type; |
| 162 | |
| 163 | jsg::Optional<jsg::Function<StartAlgorithm>> start; |
| 164 | jsg::Optional<jsg::Function<WriteAlgorithm>> write; |
| 165 | jsg::Optional<jsg::Function<AbortAlgorithm>> abort; |
| 166 | jsg::Optional<jsg::Function<CloseAlgorithm>> close; |
| 167 | |
| 168 | JSG_STRUCT(type, start, write, abort, close); |
| 169 | |
| 170 | // TODO(cleanup): Get rid of this override and parse the type directly in param-extractor.rs |
| 171 | JSG_STRUCT_TS_OVERRIDE(<W = any> { |
| 172 | write?: (chunk: W, controller: WritableStreamDefaultController) => void | Promise<void>; |
| 173 | start?: (controller: WritableStreamDefaultController) => void | Promise<void>; |
| 174 | abort?: (reason: any) => void | Promise<void>; |
| 175 | close?: () => void | Promise<void>; |
| 176 | }); |
| 177 | }; |
| 178 | |
| 179 | struct Transformer { |
| 180 | using Controller = jsg::Ref<TransformStreamDefaultController>; |
| 181 | using StartAlgorithm = jsg::Promise<void>(Controller); |
| 182 | using TransformAlgorithm = jsg::Promise<void>(v8::Local<v8::Value>, Controller); |
| 183 | using FlushAlgorithm = jsg::Promise<void>(Controller); |
| 184 | using CancelAlgorithm = jsg::Promise<void>(jsg::JsValue reason); |
| 185 | |
| 186 | jsg::Optional<kj::String> readableType; |
| 187 | jsg::Optional<kj::String> writableType; |
| 188 | |
| 189 | jsg::Optional<jsg::Function<StartAlgorithm>> start; |
| 190 | jsg::Optional<jsg::Function<TransformAlgorithm>> transform; |
| 191 | jsg::Optional<jsg::Function<FlushAlgorithm>> flush; |
| 192 | jsg::Optional<jsg::Function<CancelAlgorithm>> cancel; |
| 193 | |
| 194 | // The expectedLength is a non-standard extension used to support specifying the |
| 195 | // content-length when using a TransformStream readable side as the body of a |
| 196 | // request or response. |
| 197 | jsg::Optional<uint64_t> expectedLength; |
| 198 | |
| 199 | JSG_STRUCT(readableType, writableType, start, transform, flush, cancel, expectedLength); |
| 200 | JSG_STRUCT_TS_OVERRIDE(<I = any, O = any> { |
| 201 | start?: (controller: TransformStreamDefaultController<O>) => void | Promise<void>; |
| 202 | transform?: (chunk: I, controller: TransformStreamDefaultController<O>) => void | Promise<void>; |
| 203 | flush?: (controller: TransformStreamDefaultController<O>) => void | Promise<void>; |
| 204 | cancel?: (reason: any) => void | Promise<void>; |
| 205 | expectedLength?: number; |
| 206 | }); |
| 207 | }; |
| 208 | |
| 209 | // ReadableStreamSource and WritableStreamSink |
| 210 | // |
| 211 | // These are implementation interfaces for ReadableStream and WritableStream. If you just need to |
| 212 | // use a ReadableStream or WritableStream, you can safely skip reading this. If you need to |
| 213 | // implement a new kind of stream, read on. |
| 214 | |
| 215 | // In the original Workers streams implementation, a ReadableStream would have a |
| 216 | // ReadableStreamSource backing it. Likewise, a WritableStream would have a WritableStreamSink. |
| 217 | // The ReadableStreamSource and WritableStreamSink are kj heap objects that provide a thin |
| 218 | // wrapper on internal native stream sources originating from within the Workers runtime. |
| 219 | // |
| 220 | // With implementation of full streams standard support, we introduce the new abstraction APIs |
| 221 | // ReadableStreamController and WritableStreamController, which will provide the underlying |
| 222 | // implementation for both ReadableStream and WritableStream, respectively. |
| 223 | // |
| 224 | // When creating a new kind of *internal* ReadableStream, where the data is originating internally |
| 225 | // from a kj stream, you will still implement the ReadableStreamSource API, just as before. |
| 226 | // Likewise, when creating a new kind of *internal* WritableStream, where the data destination is |
| 227 | // a kj stream, you will implement the WritableStreamSink API. |
| 228 | |
| 229 | class WritableStreamSink { |
| 230 | public: |
| 231 | virtual kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) KJ_WARN_UNUSED_RESULT = 0; |
| 232 | virtual kj::Promise<void> write( |
| 233 | kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) KJ_WARN_UNUSED_RESULT = 0; |
| 234 | |
| 235 | virtual kj::Promise<void> end() KJ_WARN_UNUSED_RESULT = 0; |
| 236 | // Must call to flush and finish the stream. |
| 237 | |
| 238 | virtual kj::Maybe<kj::Promise<DeferredProxy<void>>> tryPumpFrom( |
| 239 | ReadableStreamSource& input, bool end); |
| 240 | |
| 241 | virtual void abort(kj::Exception reason) = 0; |
| 242 | // TODO(conform): abort() should return a promise after which closed fulfillers should be |
| 243 | // rejected. This may necessitate an "erroring" state. |
| 244 | |
| 245 | // Tells the sink that it is no longer to be responsible for encoding in the correct format. |
| 246 | // Instead, the caller takes responsibility. The expected encoding is returned; the caller |
| 247 | // promises that all future writes will use this encoding. The default implementation returns |
| 248 | // IDENTITY, which is always correct since that's the encoding write()s should have used if |
| 249 | // this weren't called at all. |
| 250 | virtual StreamEncoding disownEncodingResponsibility() { |
| 251 | return StreamEncoding::IDENTITY; |
| 252 | } |
| 253 | }; |
| 254 | |
| 255 | class ReadableStreamSource { |
| 256 | public: |
| 257 | virtual kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) = 0; |
| 258 | |
| 259 | // The ReadableStreamSource version of pumpTo() has no `amount` parameter, since the Streams spec |
| 260 | // only defines pumping everything. |
| 261 | // |
| 262 | // If `end` is true, then `output.end()` will be called after pumping. Note that it's especially |
| 263 | // important to take advantage of this when using deferred proxying since calling `end()` |
| 264 | // directly might attempt to use the `IoContext` to call `registerPendingEvent()`. |
| 265 | virtual kj::Promise<DeferredProxy<void>> pumpTo(WritableStreamSink& output, bool end); |
| 266 | |
| 267 | // If pumpTo() pumps to a system stream, what is the best encoding for that system stream to |
| 268 | // use? This is just a hint. |
| 269 | virtual StreamEncoding getPreferredEncoding() { |
| 270 | return StreamEncoding::IDENTITY; |
| 271 | }; |
| 272 | |
| 273 | virtual kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding); |
| 274 | |
| 275 | kj::Promise<kj::Array<byte>> readAllBytes(uint64_t limit); |
| 276 | kj::Promise<kj::String> readAllText( |
| 277 | uint64_t limit, ReadAllTextOption option = ReadAllTextOption::NULL_TERMINATE); |
| 278 | |
| 279 | // Hook to inform this ReadableStreamSource that the ReadableStream has been canceled. This only |
| 280 | // really means anything to TransformStreams, which are supposed to propagate the error to the |
| 281 | // writable side, and custom ReadableStreams, which we don't implement yet. |
| 282 | // |
| 283 | // NOTE: By "propagate the error back to the writable stream", I mean: if the WritableStream is in |
| 284 | // the Writable state, set it to the Errored state and reject its closed fulfiller with |
| 285 | // `reason`. I'm not sure how I'm going to do this yet. |
| 286 | virtual void cancel(kj::Exception reason); |
| 287 | // TODO(conform): Should return promise. |
| 288 | // |
| 289 | // TODO(conform): `reason` should be allowed to be any JS value, and not just an exception. |
| 290 | // That is, something silly like `stream.cancel(42)` should be allowed and trigger a |
| 291 | // rejection with the integer `42`. |
| 292 | |
| 293 | struct Tee { |
| 294 | kj::Own<ReadableStreamSource> branches[2]; |
| 295 | }; |
| 296 | |
| 297 | // Implement this if your ReadableStreamSource has a better way to tee a stream than the naive |
| 298 | // method, which relies upon `tryRead()`. The default implementation returns nullptr. |
| 299 | virtual kj::Maybe<Tee> tryTee(uint64_t limit); |
| 300 | }; |
| 301 | |
| 302 | struct PipeToOptions { |
| 303 | jsg::Optional<bool> preventAbort; |
| 304 | jsg::Optional<bool> preventCancel; |
| 305 | jsg::Optional<bool> preventClose; |
| 306 | jsg::Optional<jsg::Ref<AbortSignal>> signal; |
| 307 | |
| 308 | JSG_STRUCT(preventAbort, preventCancel, preventClose, signal); |
| 309 | JSG_STRUCT_TS_OVERRIDE(StreamPipeOptions); |
| 310 | |
| 311 | // An additional, internal only property that is used to indicate |
| 312 | // when the pipe operation is used for a pipeThrough rather than |
| 313 | // a pipeTo. We use this information, for instance, to identify |
| 314 | // when we should mark returned rejected promises as handled. |
| 315 | bool pipeThrough = false; |
| 316 | }; |
| 317 | |
| 318 | namespace StreamStates { |
| 319 | struct Closed { |
| 320 | static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj; |
| 321 | }; |
| 322 | using Errored = jsg::Value; |
| 323 | struct Erroring { |
| 324 | static constexpr kj::StringPtr NAME KJ_UNUSED = "erroring"_kj; |
| 325 | jsg::Value reason; |
| 326 | |
| 327 | Erroring(jsg::Value reason): reason(kj::mv(reason)) {} |
| 328 | |
| 329 | void visitForGc(jsg::GcVisitor& visitor) { |
| 330 | visitor.visit(reason); |
| 331 | } |
| 332 | }; |
| 333 | } // namespace StreamStates |
| 334 | |
| 335 | // A ReadableStreamController provides the underlying implementation for a ReadableStream. |
| 336 | // We will generally have three implementations: |
| 337 | // * ReadableStreamDefaultController |
| 338 | // * ReadableByteStreamController |
| 339 | // * ReadableStreamInternalController |
| 340 | // |
| 341 | // The ReadableStreamDefaultController and ReadableByteStreamController are defined by the |
| 342 | // streams standard and source all of the stream data from JavaScript functions provided by |
| 343 | // user code. |
| 344 | // |
| 345 | // The ReadableStreamInternalController is Workers runtime specific and provides a bridge |
| 346 | // to the existing ReadableStreamSource API. At the API contract layer, the |
| 347 | // ReadableByteStreamController and ReadableStreamInternalController will appear to be |
| 348 | // identical. Internally, however, they will be very different from one another. |
| 349 | // |
| 350 | // The ReadableStreamController instance is meant to be a private member of the ReadableStream, |
| 351 | // e.g. |
| 352 | // class ReadableStream { |
| 353 | // public: |
| 354 | // // ... |
| 355 | // private: |
| 356 | // ReadableStreamController controller; |
| 357 | // // ... |
| 358 | // } |
| 359 | // |
| 360 | // As such, it exists within the V8 heap (it's allocated directly as a member of the |
| 361 | // ReadableStream) and will always execute within the V8 isolate lock. |
| 362 | // |
| 363 | // The methods here return jsg::Promise rather than kj::Promise because the controller |
| 364 | // operations here do not always require passing through the kj mechanisms or kj event loop. |
| 365 | // Likewise, we do not make use of kj::Exception in these interfaces because the stream |
| 366 | // standard dictates that streams can be canceled/aborted/errored using any arbitrary JavaScript |
| 367 | // value, not just Errors. |
| 368 | class ReadableStreamController { |
| 369 | public: |
| 370 | // The ReadableStreamController::Reader interface is a base for all ReadableStream reader |
| 371 | // implementations and is used solely as a means of attaching a Reader implementation to |
| 372 | // the internal state of the controller. See the ReadableStream::*Reader classes for the |
| 373 | // full Reader API. |
| 374 | class Reader { |
| 375 | public: |
| 376 | // True if the reader is a BYOB reader. |
| 377 | virtual bool isByteOriented() const = 0; |
| 378 | |
| 379 | // When a Reader is locked to a controller, the controller will attach itself to the reader, |
| 380 | // passing along the closed promise that will be used to communicate state to the |
| 381 | // user code. |
| 382 | // |
| 383 | // The Reader will hold a reference to the controller that will be cleared when the reader |
| 384 | // is released or destroyed. The controller is guaranteed to either outlive or detach the |
| 385 | // reader so the ReadableStreamController& reference should remain valid. |
| 386 | virtual void attach(ReadableStreamController& controller, jsg::Promise<void> closedPromise) = 0; |
| 387 | |
| 388 | // When a Reader lock is released, the controller will signal to the reader that it has been |
| 389 | // detached. |
| 390 | virtual void detach() = 0; |
| 391 | }; |
| 392 | |
| 393 | struct ByobOptions { |
| 394 | static constexpr size_t DEFAULT_AT_LEAST = 1; |
| 395 | |
| 396 | jsg::V8Ref<v8::ArrayBufferView> bufferView; |
| 397 | size_t byteOffset = 0; |
| 398 | size_t byteLength; |
| 399 | |
| 400 | // The minimum number of elements that should be read. When not specified, the default |
| 401 | // is DEFAULT_AT_LEAST. This is a non-standard, Workers-specific extension to |
| 402 | // support the readAtLeast method on the ReadableStreamBYOBReader object. |
| 403 | // ReaderImpl::read() converts this to bytes by multiplying by element size. |
| 404 | kj::Maybe<size_t> atLeast = DEFAULT_AT_LEAST; |
| 405 | |
| 406 | // True if the given buffer should be detached. Per the spec, we should always be |
| 407 | // detaching a BYOB buffer but the original Workers implementation did not. |
| 408 | // To avoid breaking backwards compatibility, a compatibility flag is provided to turn |
| 409 | // detach on/off as appropriate. |
| 410 | bool detachBuffer = true; |
| 411 | }; |
| 412 | |
| 413 | struct Tee { |
| 414 | jsg::Ref<ReadableStream> branch1; |
| 415 | jsg::Ref<ReadableStream> branch2; |
| 416 | }; |
| 417 | |
| 418 | // Abstract API for ReadableStreamController implementations that provide their own |
| 419 | // tee implementations that are not backed by kj's tee. Each branch of the tee uses |
| 420 | // the TeeController to interface with the shared underlying source, and the |
| 421 | // TeeController ensures that each Branch receives the data that is read. |
| 422 | class TeeController { |
| 423 | public: |
| 424 | // Represents an individual ReadableStreamController tee branch registered with |
| 425 | // a TeeController. One or more branches is registered with the TeeController. |
| 426 | class Branch { |
| 427 | public: |
| 428 | virtual ~Branch() noexcept(false) {} |
| 429 | |
| 430 | virtual void doClose(jsg::Lock& js) = 0; |
| 431 | virtual void doError(jsg::Lock& js, v8::Local<v8::Value> reason) = 0; |
| 432 | virtual void handleData(jsg::Lock& js, ReadResult result) = 0; |
| 433 | }; |
| 434 | |
| 435 | class BranchPtr { |
| 436 | public: |
| 437 | inline BranchPtr(Branch* branch): inner(branch) { |
| 438 | KJ_ASSERT(inner != nullptr); |
| 439 | } |
| 440 | BranchPtr(BranchPtr&& other) = default; |
| 441 | BranchPtr& operator=(BranchPtr&&) = default; |
| 442 | BranchPtr(BranchPtr& other) = default; |
| 443 | |
| 444 | inline void doClose(jsg::Lock& js) { |
| 445 | inner->doClose(js); |
| 446 | } |
| 447 | |
| 448 | inline void doError(jsg::Lock& js, v8::Local<v8::Value> reason) { |
| 449 | inner->doError(js, reason); |
| 450 | } |
| 451 | |
| 452 | inline void handleData(jsg::Lock& js, ReadResult result) { |
| 453 | inner->handleData(js, kj::mv(result)); |
| 454 | } |
| 455 | |
| 456 | inline uint hashCode() { |
| 457 | return kj::hashCode(inner); |
| 458 | } |
| 459 | inline bool operator==(BranchPtr& other) const { |
| 460 | return inner == other.inner; |
| 461 | } |
| 462 | |
| 463 | private: |
| 464 | Branch* inner; |
| 465 | }; |
| 466 | |
| 467 | virtual ~TeeController() noexcept(false) {} |
| 468 | |
| 469 | virtual void addBranch(Branch* branch) = 0; |
| 470 | |
| 471 | virtual void close(jsg::Lock& js) = 0; |
| 472 | |
| 473 | virtual void error(jsg::Lock& js, v8::Local<v8::Value> reason) = 0; |
| 474 | |
| 475 | virtual void ensurePulling(jsg::Lock& js) = 0; |
| 476 | |
| 477 | // maybeJs will be nullptr when the isolate lock is not available. |
| 478 | // If maybeJs is set, any operations pending for the branch will be canceled. |
| 479 | virtual void removeBranch(Branch* branch, kj::Maybe<jsg::Lock&> maybeJs) = 0; |
| 480 | }; |
| 481 | |
| 482 | // The PipeController simplifies the abstraction between ReadableStreamController |
| 483 | // and WritableStreamController so that the pipeTo/pipeThrough/tryPipeTo can work |
| 484 | // without caring about what kind of controller it is working with. |
| 485 | class PipeController { |
| 486 | public: |
| 487 | virtual ~PipeController() noexcept(false) {} |
| 488 | virtual bool isClosed() = 0; |
| 489 | virtual kj::Maybe<v8::Local<v8::Value>> tryGetErrored(jsg::Lock& js) = 0; |
| 490 | virtual void cancel(jsg::Lock& js, v8::Local<v8::Value> reason) = 0; |
| 491 | virtual void close(jsg::Lock& js) = 0; |
| 492 | virtual void error(jsg::Lock& js, v8::Local<v8::Value> reason) = 0; |
| 493 | virtual void release(jsg::Lock& js, kj::Maybe<v8::Local<v8::Value>> maybeError = kj::none) = 0; |
| 494 | virtual kj::Maybe<kj::Promise<void>> tryPumpTo(WritableStreamSink& sink, bool end) = 0; |
| 495 | virtual jsg::Promise<ReadResult> read(jsg::Lock& js) = 0; |
| 496 | }; |
| 497 | |
| 498 | virtual ~ReadableStreamController() noexcept(false) {} |
| 499 | |
| 500 | virtual void setOwnerRef(ReadableStream& stream) = 0; |
| 501 | |
| 502 | virtual jsg::Ref<ReadableStream> addRef() = 0; |
| 503 | |
| 504 | // Returns true if the underlying source for this controller is byte-oriented and |
| 505 | // therefore supports the pull into API. When false, the stream can be used to pass |
| 506 | // any arbitrary JavaScript value through. |
| 507 | virtual bool isByteOriented() const = 0; |
| 508 | |
| 509 | // Reads data from the stream. If the stream is byte-oriented, then the ByobOptions can be |
| 510 | // specified to provide a v8::ArrayBuffer to be filled by the read operation. If the ByobOptions |
| 511 | // are provided and the stream is not byte-oriented, the operation will return a rejected promise. |
| 512 | virtual kj::Maybe<jsg::Promise<ReadResult>> read( |
| 513 | jsg::Lock& js, kj::Maybe<ByobOptions> byobOptions) = 0; |
| 514 | |
| 515 | // Performs a draining read operation that: |
| 516 | // 1. Drains all currently buffered data from the queue |
| 517 | // 2. Pumps the controller for synchronously available data (respecting pull promise state) |
| 518 | // 3. Returns bytes even for value streams (converting ArrayBuffer/ArrayBufferView/string) |
| 519 | // 4. Has mutual exclusion with regular reads - returns rejected promise if pending regular reads |
| 520 | // 5. Returns done: true with final data when stream is closing |
| 521 | // |
| 522 | // This is a C++ only API (not exposed to JavaScript) intended for optimized pipe operations. |
| 523 | // Returns kj::none if the stream is locked in a way that prevents the read. |
| 524 | // |
| 525 | // The maxRead parameter provides a soft limit on how much data to read. Both the initial |
| 526 | // buffer drain and subsequent synchronous pump attempts stop when the total bytes read |
| 527 | // reaches maxRead (after finishing the current item). This prevents unbounded memory |
| 528 | // accumulation when a fast producer outpaces a slow consumer. |
| 529 | virtual kj::Maybe<jsg::Promise<DrainingReadResult>> drainingRead( |
| 530 | jsg::Lock& js, size_t maxRead = kj::maxValue) = 0; |
| 531 | |
| 532 | // The pipeTo implementation fully consumes the stream by directing all of its data at the |
| 533 | // destination. Controllers should try to be as efficient as possible here. For instance, if |
| 534 | // a ReadableStreamInternalController is piping to a WritableStreamInternalController, then |
| 535 | // a more efficient kj pipe should be possible. |
| 536 | virtual jsg::Promise<void> pipeTo( |
| 537 | jsg::Lock& js, WritableStreamController& destination, PipeToOptions options) = 0; |
| 538 | |
| 539 | // Indicates that the consumer no longer has any interest in the streams data. |
| 540 | virtual jsg::Promise<void> cancel(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) = 0; |
| 541 | |
| 542 | // Branches the ReadableStreamController into two ReadableStream instances that will receive |
| 543 | // this streams data. The specific details of how the branching occurs is entirely up to the |
| 544 | // controller implementation. |
| 545 | virtual Tee tee(jsg::Lock& js) = 0; |
| 546 | |
| 547 | virtual bool isClosedOrErrored() const = 0; |
| 548 | |
| 549 | virtual bool isClosed() const = 0; |
| 550 | |
| 551 | virtual bool isDisturbed() = 0; |
| 552 | |
| 553 | // True if a Reader has been locked to this controller. |
| 554 | virtual bool isLockedToReader() const = 0; |
| 555 | |
| 556 | // Locks this controller to the given reader, returning true if the lock was successful, or false |
| 557 | // if the controller was already locked. |
| 558 | virtual bool lockReader(jsg::Lock& js, Reader& reader) = 0; |
| 559 | |
| 560 | // Removes the lock and releases the reader from this controller. |
| 561 | // maybeJs will be nullptr when the isolate lock is not available. |
| 562 | // If maybeJs is set, the reader's closed promise will be resolved. |
| 563 | virtual void releaseReader(Reader& reader, kj::Maybe<jsg::Lock&> maybeJs) = 0; |
| 564 | |
| 565 | virtual kj::Maybe<PipeController&> tryPipeLock() = 0; |
| 566 | |
| 567 | virtual void visitForGc(jsg::GcVisitor& visitor) {}; |
| 568 | |
| 569 | // Fully consumes the ReadableStream. If the stream is already locked to a reader or |
| 570 | // errored, the returned JS promise will reject. If the stream is already closed, the |
| 571 | // returned JS promise will resolve with a zero-length result. Importantly, this will |
| 572 | // lock the stream and will fully consume it. |
| 573 | // |
| 574 | // limit specifies an upper maximum bound on the number of bytes permitted to be read. |
| 575 | // The promise will reject if the read will produce more bytes than the limit. |
| 576 | virtual jsg::Promise<jsg::BufferSource> readAllBytes(jsg::Lock& js, uint64_t limit) = 0; |
| 577 | |
| 578 | // Fully consumes the ReadableStream. If the stream is already locked to a reader or |
| 579 | // errored, the returned JS promise will reject. If the stream is already closed, the |
| 580 | // returned JS promise will resolve with a zero-length result. Importantly, this will |
| 581 | // lock the stream and will fully consume it. |
| 582 | // |
| 583 | // limit specifies an upper maximum bound on the number of bytes permitted to be read. |
| 584 | // The promise will reject if the read will produce more bytes than the limit. |
| 585 | virtual jsg::Promise<kj::String> readAllText(jsg::Lock& js, uint64_t limit) = 0; |
| 586 | |
| 587 | virtual kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding) = 0; |
| 588 | |
| 589 | virtual void setup(jsg::Lock& js, |
| 590 | jsg::Optional<UnderlyingSource> maybeUnderlyingSource, |
| 591 | jsg::Optional<StreamQueuingStrategy> maybeQueuingStrategy) {} |
| 592 | |
| 593 | virtual kj::Promise<DeferredProxy<void>> pumpTo( |
| 594 | jsg::Lock& js, kj::Own<WritableStreamSink> sink, bool end) = 0; |
| 595 | |
| 596 | // If pumpTo() pumps to a system stream, what is the best encoding for that system stream to |
| 597 | // use? This is just a hint. |
| 598 | virtual StreamEncoding getPreferredEncoding() { |
| 599 | return StreamEncoding::IDENTITY; |
| 600 | } |
| 601 | |
| 602 | virtual kj::Own<ReadableStreamController> detach(jsg::Lock& js, bool ignoreDisturbed) = 0; |
| 603 | |
| 604 | // Used by sockets to signal that the ReadableStream shouldn't allow reads due to pending |
| 605 | // closure. |
| 606 | virtual void setPendingClosure() = 0; |
| 607 | |
| 608 | virtual kj::StringPtr jsgGetMemoryName() const = 0; |
| 609 | virtual size_t jsgGetMemorySelfSize() const = 0; |
| 610 | virtual void jsgGetMemoryInfo(jsg::MemoryTracker& tracker) const = 0; |
| 611 | }; |
| 612 | |
| 613 | kj::Own<ReadableStreamController> newReadableStreamJsController(); |
| 614 | kj::Own<ReadableStreamController> newReadableStreamInternalController( |
| 615 | IoContext& ioContext, kj::Own<ReadableStreamSource> source); |
| 616 | |
| 617 | // A WritableStreamController provides the underlying implementation for a WritableStream. |
| 618 | // We will generally have two implementations: |
| 619 | // * WritableStreamDefaultController |
| 620 | // * WritableStreamInternalController |
| 621 | // |
| 622 | // The WritableStreamDefaultController is defined by the streams standard and directs all |
| 623 | // of the stream data to JavaScript functions provided by user code. |
| 624 | // |
| 625 | // The WritableStreamInternalController is Workers runtime specific and provides a bridge |
| 626 | // to the existing WritableStreamSink API. |
| 627 | // |
| 628 | // The WritableStreamController instance is meant to be a private member of the WritableStream, |
| 629 | // e.g. |
| 630 | // class WritableStream { |
| 631 | // public: |
| 632 | // // ... |
| 633 | // private: |
| 634 | // WritableStreamController controller; |
| 635 | // }; |
| 636 | // |
| 637 | // As such, it exists within the V8 heap (it's allocated directly as a member of the |
| 638 | // WritableStream) and will always execute within the V8 isolate lock. |
| 639 | // |
| 640 | // The methods here return jsg::Promise rather than kj::Promise because the controller |
| 641 | // operations here do not always require passing through the kj mechanisms or kj event loop. |
| 642 | // Likewise, we do not make use of kj::Exception in these interfaces because the stream |
| 643 | // standard dictates that streams can be canceled/aborted/errored using any arbitrary JavaScript |
| 644 | // value, not just Errors. |
| 645 | class WritableStreamController { |
| 646 | public: |
| 647 | // The WritableStreamController::Writer interface is a base for all WritableStream writer |
| 648 | // implementations and is used solely as a means of attaching a Writer implementation to |
| 649 | // the internal state of the controller. See the WritableStream::*Writer classes for the |
| 650 | // full Writer API. |
| 651 | class Writer { |
| 652 | public: |
| 653 | // When a Writer is locked to a controller, the controller will attach itself to the writer, |
| 654 | // passing along the closed and ready promises that will be used to communicate state to the |
| 655 | // user code. |
| 656 | // |
| 657 | // The controller is guaranteed to either outlive the Writer or will detach the Writer so the |
| 658 | // WritableStreamController& reference should always remain valid. |
| 659 | virtual void attach(jsg::Lock& js, |
| 660 | WritableStreamController& controller, |
| 661 | jsg::Promise<void> closedPromise, |
| 662 | jsg::Promise<void> readyPromise) = 0; |
| 663 | |
| 664 | // When a Writer lock is released, the controller will signal to the writer that is has been |
| 665 | // detached. |
| 666 | virtual void detach() = 0; |
| 667 | |
| 668 | // The ready promise can be replaced whenever backpressure is signaled by the underlying |
| 669 | // controller. |
| 670 | virtual void replaceReadyPromise(jsg::Lock& js, jsg::Promise<void> readyPromise) = 0; |
| 671 | }; |
| 672 | |
| 673 | struct PendingAbort { |
| 674 | kj::Maybe<jsg::Promise<void>::Resolver> resolver; |
| 675 | jsg::Promise<void> promise; |
| 676 | jsg::Value reason; |
| 677 | bool reject = false; |
| 678 | |
| 679 | PendingAbort(jsg::Lock& js, |
| 680 | jsg::PromiseResolverPair<void> prp, |
| 681 | v8::Local<v8::Value> reason, |
| 682 | bool reject); |
| 683 | |
| 684 | PendingAbort(jsg::Lock& js, v8::Local<v8::Value> reason, bool reject); |
| 685 | |
| 686 | void complete(jsg::Lock& js); |
| 687 | |
| 688 | void fail(jsg::Lock& js, v8::Local<v8::Value> reason); |
| 689 | |
| 690 | inline jsg::Promise<void> whenResolved(jsg::Lock& js) { |
| 691 | return promise.whenResolved(js); |
| 692 | } |
| 693 | |
| 694 | inline jsg::Promise<void> whenResolved(auto&& func) { |
| 695 | return promise.whenResolved(kj::fwd(func)); |
| 696 | } |
| 697 | |
| 698 | inline jsg::Promise<void> whenResolved(auto&& func, auto&& errFunc) { |
| 699 | return promise.whenResolved(kj::fwd(func), kj::fwd(errFunc)); |
| 700 | } |
| 701 | |
| 702 | void visitForGc(jsg::GcVisitor& visitor) { |
| 703 | visitor.visit(resolver, promise, reason); |
| 704 | } |
| 705 | |
| 706 | static kj::Maybe<kj::Own<PendingAbort>> dequeue( |
| 707 | kj::Maybe<kj::Own<PendingAbort>>& maybePendingAbort); |
| 708 | |
| 709 | JSG_MEMORY_INFO(PendingAbort) { |
| 710 | tracker.trackField("resolver", resolver); |
| 711 | tracker.trackField("promise", promise); |
| 712 | tracker.trackField("reason", reason); |
| 713 | } |
| 714 | }; |
| 715 | |
| 716 | virtual ~WritableStreamController() noexcept(false) {} |
| 717 | |
| 718 | virtual void setOwnerRef(WritableStream& stream) = 0; |
| 719 | |
| 720 | virtual jsg::Ref<WritableStream> addRef() = 0; |
| 721 | |
| 722 | // The controller implementation will determine what kind of JavaScript data |
| 723 | // it is capable of writing, returning a rejected promise if the written |
| 724 | // data type is not supported. |
| 725 | virtual jsg::Promise<void> write(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> value) = 0; |
| 726 | |
| 727 | // Indicates that no additional data will be written to the controller. All |
| 728 | // existing pending writes should be allowed to complete. |
| 729 | virtual jsg::Promise<void> close(jsg::Lock& js, bool markAsHandled = false) = 0; |
| 730 | |
| 731 | // Waits for pending data to be written. The returned promise is resolved when all pending writes |
| 732 | // have completed. |
| 733 | virtual jsg::Promise<void> flush(jsg::Lock& js, bool markAsHandled = false) = 0; |
| 734 | |
| 735 | // Immediately interrupts existing pending writes and errors the stream. |
| 736 | virtual jsg::Promise<void> abort(jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) = 0; |
| 737 | |
| 738 | // The tryPipeFrom attempts to establish a data pipe where source's data |
| 739 | // is delivered to this WritableStreamController as efficiently as possible. |
| 740 | virtual kj::Maybe<jsg::Promise<void>> tryPipeFrom( |
| 741 | jsg::Lock& js, jsg::Ref<ReadableStream> source, PipeToOptions options) = 0; |
| 742 | |
| 743 | // Only byte-oriented WritableStreamController implementations will have a WritableStreamSink |
| 744 | // that can be detached using removeSink. A nullptr should be returned by any controller that |
| 745 | // does not support removing the sink. After the WritableStreamSink has been released, all other |
| 746 | // methods on the controller should fail with an exception as the WritableStreamSink should be |
| 747 | // the only way to interact with the underlying sink. |
| 748 | virtual kj::Maybe<kj::Own<WritableStreamSink>> removeSink(jsg::Lock& js) = 0; |
| 749 | |
| 750 | // Detaches the WritableStreamController from its underlying implementation, leaving the |
| 751 | // writable stream locked and in a state where no further writes can be made. |
| 752 | virtual void detach(jsg::Lock& js) = 0; |
| 753 | |
| 754 | virtual kj::Maybe<int> getDesiredSize() = 0; |
| 755 | |
| 756 | // True if a Writer has been locked to this controller. |
| 757 | virtual bool isLockedToWriter() const = 0; |
| 758 | |
| 759 | // Locks this controller to the given writer, returning true if the lock was successful, or false |
| 760 | // if the controller was already locked. |
| 761 | virtual bool lockWriter(jsg::Lock& js, Writer& writer) = 0; |
| 762 | |
| 763 | // Removes the lock and releases the writer from this controller. |
| 764 | // maybeJs will be nullptr when the isolate lock is not available. |
| 765 | // If maybeJs is set, the writer's closed and ready promises will be resolved. |
| 766 | virtual void releaseWriter(Writer& writer, kj::Maybe<jsg::Lock&> maybeJs) = 0; |
| 767 | |
| 768 | virtual kj::Maybe<v8::Local<v8::Value>> isErroring(jsg::Lock& js) = 0; |
| 769 | |
| 770 | virtual void visitForGc(jsg::GcVisitor& visitor) {}; |
| 771 | |
| 772 | virtual void setup(jsg::Lock& js, |
| 773 | jsg::Optional<UnderlyingSink> underlyingSink, |
| 774 | jsg::Optional<StreamQueuingStrategy> queuingStrategy) {} |
| 775 | |
| 776 | virtual bool isClosedOrClosing() = 0; |
| 777 | virtual bool isErrored() = 0; |
| 778 | |
| 779 | // True is this controller requires ArrayBuffer(Views) to be written to it. |
| 780 | virtual bool isByteOriented() const = 0; |
| 781 | |
| 782 | // Used by sockets to signal that the WritableStream shouldn't allow writes due to pending |
| 783 | // closure. |
| 784 | virtual void setPendingClosure() = 0; |
| 785 | |
| 786 | // For memory tracking |
| 787 | virtual kj::StringPtr jsgGetMemoryName() const = 0; |
| 788 | virtual size_t jsgGetMemorySelfSize() const = 0; |
| 789 | virtual void jsgGetMemoryInfo(jsg::MemoryTracker& info) const = 0; |
| 790 | }; |
| 791 | |
| 792 | kj::Own<WritableStreamController> newWritableStreamJsController(); |
| 793 | kj::Own<WritableStreamController> newWritableStreamInternalController(IoContext& ioContext, |
| 794 | kj::Own<WritableStreamSink> sink, |
| 795 | kj::Maybe<kj::Own<ByteStreamObserver>> observer, |
| 796 | kj::Maybe<uint64_t> maybeHighWaterMark = kj::none, |
| 797 | kj::Maybe<jsg::Promise<void>> maybeClosureWaitable = kj::none); |
| 798 | |
| 799 | struct Unlocked { |
| 800 | static constexpr kj::StringPtr NAME KJ_UNUSED = "unlocked"_kj; |
| 801 | }; |
| 802 | struct Locked { |
| 803 | static constexpr kj::StringPtr NAME KJ_UNUSED = "locked"_kj; |
| 804 | }; |
| 805 | |
| 806 | // When a reader is locked to a ReadableStream, a ReaderLock instance |
| 807 | // is used internally to represent the locked state in the ReadableStreamController. |
| 808 | class ReaderLocked { |
| 809 | public: |
| 810 | static constexpr kj::StringPtr NAME KJ_UNUSED = "reader-locked"_kj; |
| 811 | ReaderLocked(ReadableStreamController::Reader& reader, |
| 812 | jsg::Promise<void>::Resolver closedFulfiller, |
| 813 | kj::Maybe<IoOwn<kj::Canceler>> canceler = kj::none) |
| 814 | : reader(reader), |
| 815 | closedFulfiller(kj::mv(closedFulfiller)), |
| 816 | canceler(kj::mv(canceler)) {} |
| 817 | |
| 818 | ReaderLocked(ReaderLocked&&) = default; |
| 819 | ~ReaderLocked() noexcept(false) { |
| 820 | KJ_IF_SOME(r, reader) { |
| 821 | r.detach(); |
| 822 | } |
| 823 | } |
| 824 | KJ_DISALLOW_COPY(ReaderLocked); |
| 825 | |
| 826 | void visitForGc(jsg::GcVisitor& visitor) { |
| 827 | visitor.visit(closedFulfiller); |
| 828 | } |
| 829 | |
| 830 | ReadableStreamController::Reader& getReader() { |
| 831 | return KJ_ASSERT_NONNULL(reader); |
| 832 | } |
| 833 | |
| 834 | kj::Maybe<jsg::Promise<void>::Resolver>& getClosedFulfiller() { |
| 835 | return closedFulfiller; |
| 836 | } |
| 837 | |
| 838 | kj::Maybe<IoOwn<kj::Canceler>>& getCanceler() { |
| 839 | return canceler; |
| 840 | } |
| 841 | |
| 842 | void clear() { |
| 843 | reader = kj::none; |
| 844 | closedFulfiller = kj::none; |
| 845 | canceler = kj::none; |
| 846 | } |
| 847 | |
| 848 | JSG_MEMORY_INFO(ReaderLocked) { |
| 849 | tracker.trackField("closedFulfiller", closedFulfiller); |
| 850 | tracker.trackFieldWithSize("IoOwn<kj::Canceler>", sizeof(IoOwn<kj::Canceler>)); |
| 851 | } |
| 852 | |
| 853 | private: |
| 854 | kj::Maybe<ReadableStreamController::Reader&> reader; |
| 855 | kj::Maybe<jsg::Promise<void>::Resolver> closedFulfiller; |
| 856 | kj::Maybe<IoOwn<kj::Canceler>> canceler; |
| 857 | }; |
| 858 | |
| 859 | // When a writer is locked to a WritableStream, a WriterLock instance |
| 860 | // is used internally to represent the locked state in the WritableStreamController. |
| 861 | class WriterLocked { |
| 862 | public: |
| 863 | static constexpr kj::StringPtr NAME KJ_UNUSED = "writer-locked"_kj; |
| 864 | WriterLocked(WritableStreamController::Writer& writer, |
| 865 | jsg::Promise<void>::Resolver closedFulfiller, |
| 866 | kj::Maybe<jsg::Promise<void>::Resolver> readyFulfiller = kj::none) |
| 867 | : writer(writer), |
| 868 | closedFulfiller(kj::mv(closedFulfiller)), |
| 869 | readyFulfiller(kj::mv(readyFulfiller)) {} |
| 870 | |
| 871 | WriterLocked(WriterLocked&&) = default; |
| 872 | ~WriterLocked() noexcept(false) { |
| 873 | KJ_IF_SOME(w, writer) { |
| 874 | w.detach(); |
| 875 | } |
| 876 | } |
| 877 | |
| 878 | void visitForGc(jsg::GcVisitor& visitor) { |
| 879 | visitor.visit(closedFulfiller, readyFulfiller); |
| 880 | } |
| 881 | |
| 882 | WritableStreamController::Writer& getWriter() { |
| 883 | return KJ_ASSERT_NONNULL(writer); |
| 884 | } |
| 885 | |
| 886 | kj::Maybe<jsg::Promise<void>::Resolver>& getClosedFulfiller() { |
| 887 | return closedFulfiller; |
| 888 | } |
| 889 | |
| 890 | kj::Maybe<jsg::Promise<void>::Resolver>& getReadyFulfiller() { |
| 891 | return readyFulfiller; |
| 892 | } |
| 893 | |
| 894 | void setReadyFulfiller(jsg::Lock& js, jsg::PromiseResolverPair<void>& pair) { |
| 895 | KJ_IF_SOME(w, writer) { |
| 896 | readyFulfiller = kj::mv(pair.resolver); |
| 897 | w.replaceReadyPromise(js, kj::mv(pair.promise)); |
| 898 | } |
| 899 | } |
| 900 | |
| 901 | void clear() { |
| 902 | writer = kj::none; |
| 903 | closedFulfiller = kj::none; |
| 904 | readyFulfiller = kj::none; |
| 905 | } |
| 906 | |
| 907 | JSG_MEMORY_INFO(WriterLocked) { |
| 908 | tracker.trackField("closedFulfiller", closedFulfiller); |
| 909 | tracker.trackField("readyFulfiller", readyFulfiller); |
| 910 | } |
| 911 | |
| 912 | private: |
| 913 | kj::Maybe<WritableStreamController::Writer&> writer; |
| 914 | kj::Maybe<jsg::Promise<void>::Resolver> closedFulfiller; |
| 915 | kj::Maybe<jsg::Promise<void>::Resolver> readyFulfiller; |
| 916 | }; |
| 917 | |
| 918 | template <typename T> |
| 919 | void maybeResolvePromise( |
| 920 | jsg::Lock& js, kj::Maybe<typename jsg::Promise<T>::Resolver>& maybeResolver, T&& t) { |
| 921 | KJ_IF_SOME(resolver, maybeResolver) { |
| 922 | resolver.resolve(js, kj::fwd<T>(t)); |
| 923 | maybeResolver = nullptr; |
| 924 | } |
| 925 | } |
| 926 | |
| 927 | inline void maybeResolvePromise( |
| 928 | jsg::Lock& js, kj::Maybe<jsg::Promise<void>::Resolver>& maybeResolver) { |
| 929 | KJ_IF_SOME(resolver, maybeResolver) { |
| 930 | resolver.resolve(js); |
| 931 | maybeResolver = kj::none; |
| 932 | } |
| 933 | } |
| 934 | |
| 935 | template <typename T> |
| 936 | void maybeRejectPromise(jsg::Lock& js, |
| 937 | kj::Maybe<typename jsg::Promise<T>::Resolver>& maybeResolver, |
| 938 | v8::Local<v8::Value> reason) { |
| 939 | KJ_IF_SOME(resolver, maybeResolver) { |
| 940 | resolver.reject(js, reason); |
| 941 | maybeResolver = kj::none; |
| 942 | } |
| 943 | } |
| 944 | |
| 945 | template <typename T> |
| 946 | jsg::Promise<T> rejectedMaybeHandledPromise( |
| 947 | jsg::Lock& js, v8::Local<v8::Value> reason, bool handled) { |
| 948 | auto prp = js.newPromiseAndResolver<T>(); |
| 949 | if (handled) { |
| 950 | prp.promise.markAsHandled(js); |
| 951 | } |
| 952 | prp.resolver.reject(js, reason); |
| 953 | return kj::mv(prp.promise); |
| 954 | } |
| 955 | |
| 956 | inline kj::Maybe<IoContext&> tryGetIoContext() { |
| 957 | // TODO(cleanup): This function is obsolete; callers should just call IoContext::tryCurrent() |
| 958 | return IoContext::tryCurrent(); |
| 959 | } |
| 960 | |
| 961 | } // namespace workerd::api |