File
Blob: src/workerd/api/streams/readable-source-adapter.h
| 1 | #include "common.h" |
| 2 | #include "readable-source.h" |
| 3 | #include "readable.h" |
| 4 | |
| 5 | #include <workerd/util/state-machine.h> |
| 6 | |
| 7 | namespace workerd::api::streams { |
| 8 | |
| 9 | // We provide two utility adapters here: ReadableStreamSourceJsAdapter and |
| 10 | // ReadableSourceKjAdapter. |
| 11 | // |
| 12 | // ReadableStreamSourceJsAdapter adapts a ReadableStreamSource to a JavaScript-friendly |
| 13 | // interface. It provides methods that return JavaScript promises and use |
| 14 | // JavaScript types. It is intended to be used by JavaScript code that wants |
| 15 | // to read from a kj-backed stream. It takes ownership of the ReadableStreamSource |
| 16 | // and holds it with an IoOwn, ensures that all operations are performed on the |
| 17 | // correct IoContext, and safely cleans up after itself if the adapter is dropped. |
| 18 | // |
| 19 | // ┌───────────────────────────────────────────┐ |
| 20 | // │ ReadableStreamSourceJsAdapter │ |
| 21 | // │ │ |
| 22 | // │ ┌─────────────────────────────────────┐ │ |
| 23 | // │ │ JavaScript API │ │ |
| 24 | // │ │ │ │ |
| 25 | // │ │ • read() → Promise<ReadResult> │ │ |
| 26 | // │ │ • readAllText() → Promise<string> │ │ |
| 27 | // │ │ • readAllBytes() → Promise<bytes> │ │ |
| 28 | // │ │ • close() → Promise<void> │ │ |
| 29 | // │ │ • cancel(reason) │ │ |
| 30 | // │ │ • tryTee() → {branch1, branch2} │ │ |
| 31 | // │ └─────────────────────────────────────┘ │ |
| 32 | // │ │ │ |
| 33 | // │ ▼ │ |
| 34 | // │ ┌─────────────────────────────────────┐ │ |
| 35 | // │ │ State Management │ │ |
| 36 | // │ │ │ │ |
| 37 | // │ │ Active ──► Closed │ │ |
| 38 | // │ │ │ │ │ │ |
| 39 | // │ │ │ ▼ │ │ |
| 40 | // │ │ └─────► Canceled/Errored │ │ |
| 41 | // │ └─────────────────────────────────────┘ │ |
| 42 | // │ │ │ |
| 43 | // │ ▼ │ |
| 44 | // │ ┌─────────────────────────────────────┐ │ |
| 45 | // │ │ KJ Integration │ │ |
| 46 | // │ │ │ │ |
| 47 | // │ │ IoOwn<ReadableStreamSource> │ │ |
| 48 | // │ │ WeakRef for safe references │ │ |
| 49 | // │ │ IoContext-aware operations │ │ |
| 50 | // │ └─────────────────────────────────────┘ │ |
| 51 | // └───────────────────────────────────────────┘ |
| 52 | // │ |
| 53 | // ▼ |
| 54 | // ┌───────────────────────────────────────────┐ |
| 55 | // │ ReadableStreamSource │ |
| 56 | // │ (KJ Native Stream) │ |
| 57 | // │ │ |
| 58 | // │ • tryRead() │ |
| 59 | // │ • pumpTo() │ |
| 60 | // │ • tryGetLength() │ |
| 61 | // │ • cancel() │ |
| 62 | // └───────────────────────────────────────────┘ |
| 63 | // |
| 64 | // The ReadableSourceKjAdapter adapts a ReadableStream to a KJ-friendly |
| 65 | // ReadableStreamSource. It holds a strong reference to the ReadableStream and |
| 66 | // locks it with a ReadableStreamDefaultReader. It is intended to be used by |
| 67 | // KJ code that wants to read from a JavaScript-backed stream. It ensures that |
| 68 | // all operations are performed on the correct IoContext, and safely cleans up |
| 69 | // after itself if the adapter is dropped. |
| 70 | // |
| 71 | // ┌───────────────────────────────────────────┐ |
| 72 | // │ ReadableSourceKjAdapter │ |
| 73 | // │ │ |
| 74 | // │ ┌─────────────────────────────────────┐ │ |
| 75 | // │ │ KJ Native API │ │ |
| 76 | // │ │ │ │ |
| 77 | // │ │ • tryRead(minBytes, maxBytes) │ │ |
| 78 | // │ │ • pumpTo(sink, end) │ │ |
| 79 | // │ │ • tryGetLength(encoding) │ │ |
| 80 | // │ │ • cancel(exception) │ │ |
| 81 | // │ │ • getPreferredEncoding() │ │ |
| 82 | // │ │ • tryTee() → none (unsupported) │ │ |
| 83 | // │ └─────────────────────────────────────┘ │ |
| 84 | // │ │ │ |
| 85 | // │ ▼ │ |
| 86 | // │ ┌─────────────────────────────────────┐ │ |
| 87 | // │ │ State Management │ │ |
| 88 | // │ │ │ │ |
| 89 | // │ │ Active ──► Closed │ │ |
| 90 | // │ │ │ │ │ |
| 91 | // │ │ └─────► Canceled/Errored │ │ |
| 92 | // │ └─────────────────────────────────────┘ │ |
| 93 | // │ │ │ |
| 94 | // │ ▼ │ |
| 95 | // │ ┌─────────────────────────────────────┐ │ |
| 96 | // │ │ JavaScript Integration │ │ |
| 97 | // │ │ │ │ |
| 98 | // │ │ ReadableStreamDefaultReader │ │ |
| 99 | // │ │ WeakRef for safe references │ │ |
| 100 | // │ │ IoContext-aware JS operations │ │ |
| 101 | // │ │ Promise handling & async reads │ │ |
| 102 | // │ └─────────────────────────────────────┘ │ |
| 103 | // └───────────────────────────────────────────┘ |
| 104 | // │ |
| 105 | // ▼ |
| 106 | // ┌───────────────────────────────────────────┐ |
| 107 | // │ JavaScript ReadableStream │ |
| 108 | // │ │ |
| 109 | // │ • getReader() │ |
| 110 | // │ • read() → Promise<{value, done}> │ |
| 111 | // │ • cancel(reason) │ |
| 112 | // │ • locked, state properties │ |
| 113 | // └───────────────────────────────────────────┘ |
| 114 | |
| 115 | // Adapts a ReadableStreamSource to a JavaScript-friendly interface. |
| 116 | class ReadableStreamSourceJsAdapter final { |
| 117 | public: |
| 118 | ReadableStreamSourceJsAdapter( |
| 119 | jsg::Lock& js, IoContext& ioContext, kj::Own<ReadableSource> source); |
| 120 | KJ_DISALLOW_COPY_AND_MOVE(ReadableStreamSourceJsAdapter); |
| 121 | ~ReadableStreamSourceJsAdapter() noexcept(false); |
| 122 | |
| 123 | // Returns true if the adapter is closed or canceled. |
| 124 | bool isClosed(); |
| 125 | |
| 126 | // If the adapter is canceled, returns the exception it was |
| 127 | // canceled with. Otherwise returns null. |
| 128 | kj::Maybe<const kj::Exception&> isCanceled() KJ_LIFETIMEBOUND; |
| 129 | |
| 130 | // Cancels the underlying source if it is still active. If an |
| 131 | // exception is provided, the source will be errored with that. |
| 132 | // If no exception is provided, the source will be closed without |
| 133 | // error. All in-flight and pending read requests will be rejected. |
| 134 | // Unlike close(), the effect is immediate. |
| 135 | void cancel(kj::Exception exception); |
| 136 | |
| 137 | // Like cancel() but with the error reason provided as a JS value. |
| 138 | void cancel(jsg::Lock& js, const jsg::JsValue& reason); |
| 139 | |
| 140 | // Closes the stream immediatey without error if it is still |
| 141 | // active. All in-flight and pending read requests will be |
| 142 | // rejected with a cancelation error but the adapter will |
| 143 | // transition to the closed state rather than the errored state. |
| 144 | // If the adapter is already closed or canceled, this is a no-op. |
| 145 | void shutdown(jsg::Lock& js); |
| 146 | |
| 147 | // Causes the adapter to enter the closing state. Any pending |
| 148 | // read requests will be allowed to complete but no new read requests |
| 149 | // will be accepted. The underlying source will be closed fully |
| 150 | // once all pending reads complete. If cancel() has already been |
| 151 | // called, this is a no-op and the returned promise resolves |
| 152 | // immediately. If cancel() is called while this is pending, |
| 153 | // the returned promise will reject with the same exception |
| 154 | // and the cancel will supersede the close. |
| 155 | jsg::Promise<void> close(jsg::Lock& js); |
| 156 | |
| 157 | struct ReadOptions { |
| 158 | // The buffer to read into. The maximum number of bytes read |
| 159 | // is equal to the length of this buffer. The actual number of |
| 160 | // bytes read is indicated by the resolved value of the promise |
| 161 | // but will never exceed the length of this buffer. |
| 162 | jsg::BufferSource buffer; |
| 163 | |
| 164 | // The optional minimum number of bytes to read. If not provided, |
| 165 | // the read will complete as soon as at least the mininum number |
| 166 | // of bytes to satisfy the minimum bytes-per-element of the input |
| 167 | // buffer is available. |
| 168 | // It is often more efficient to provide a minimum number of bytes |
| 169 | // because it allows to implementation to wait until larger chunks |
| 170 | // of data are available before completing the read. |
| 171 | kj::Maybe<size_t> minBytes; |
| 172 | }; |
| 173 | struct ReadResult { |
| 174 | // The buffer containing the data that was read. The length |
| 175 | // of the buffer may be less than the length of the buffer |
| 176 | // provided in ReadOptions if fewer bytes were available. |
| 177 | // The identity of the underlying ArrayBuffer will be the same |
| 178 | // but the buffer itself will be a new type array view. |
| 179 | // of the same type as that provided in ReadOptions. |
| 180 | // If the read produced no data because the stream is |
| 181 | // closed, the type array will be zero length. |
| 182 | jsg::BufferSource buffer; |
| 183 | |
| 184 | // True if the stream is now closed and no further reads |
| 185 | // are possible. If this is true, the buffer will be zero |
| 186 | // length. |
| 187 | bool done = false; |
| 188 | }; |
| 189 | |
| 190 | // Submit a read request. The returned promise resolves with a |
| 191 | // BufferSource containing the data that was read. |
| 192 | jsg::Promise<ReadResult> read(jsg::Lock& js, ReadOptions options); |
| 193 | |
| 194 | // Utility function to read the entire stream as text. This is |
| 195 | // terminal in that once this is called, no further reads |
| 196 | // are possible. The entire stream will be read and concatenated |
| 197 | // and the resulting string returned. If the stream errors while |
| 198 | // reading, the promise will reject with the error. |
| 199 | // If there are pending reads when this is called, those reads |
| 200 | // will be allowed to complete first, and then the stream will |
| 201 | // be read to the end. |
| 202 | jsg::Promise<jsg::JsRef<jsg::JsString>> readAllText(jsg::Lock& js, uint64_t limit = kj::maxValue); |
| 203 | |
| 204 | // Utility function to read the entire stream as bytes. This is |
| 205 | // terminal in that once this is called, no further reads |
| 206 | // are possible. The entire stream will be read and concatenated |
| 207 | // and the resulting bytes returned as a single BufferSource. |
| 208 | // If the stream errors while reading, the promise will reject |
| 209 | // with the error. |
| 210 | // If there are pending reads when this is called, those reads |
| 211 | // will be allowed to complete first, and then the stream will |
| 212 | // be read to the end. |
| 213 | jsg::Promise<jsg::BufferSource> readAllBytes(jsg::Lock& js, uint64_t limit = kj::maxValue); |
| 214 | |
| 215 | // If the stream is still active, tries to get the total length, |
| 216 | // if known. If the length is not known, the encoding does not |
| 217 | // match the encoding of the underlying stream, or the stream is |
| 218 | // closed or errored, returns kj::none. |
| 219 | kj::Maybe<uint64_t> tryGetLength(StreamEncoding encoding); |
| 220 | |
| 221 | struct Tee { |
| 222 | kj::Own<ReadableStreamSourceJsAdapter> branch1; |
| 223 | kj::Own<ReadableStreamSourceJsAdapter> branch2; |
| 224 | }; |
| 225 | // Tees the stream into two branches. The returned Tee contains |
| 226 | // two new ReadableStreamSourceJsAdapter instances that will |
| 227 | // each receive the same data as this instance. Once this is called, |
| 228 | // this instance is no longer usable and all further operations |
| 229 | // on it will fail. Each branch operates independently; closing, |
| 230 | // canceling, or erroring one branch has no effect on the other branch. |
| 231 | // If this instance is already closed or canceled, or if there are |
| 232 | // in-flight or pending reads, this will throw. |
| 233 | kj::Maybe<Tee> tryTee(jsg::Lock& js, uint64_t limit = kj::maxValue); |
| 234 | |
| 235 | private: |
| 236 | struct Active; |
| 237 | struct Closed final { |
| 238 | static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj; |
| 239 | }; |
| 240 | struct Open { |
| 241 | static constexpr kj::StringPtr NAME KJ_UNUSED = "open"_kj; |
| 242 | IoOwn<Active> active; |
| 243 | }; |
| 244 | |
| 245 | // State machine for tracking readable source adapter lifecycle: |
| 246 | // Open -> Closed (normal close) |
| 247 | // Open -> kj::Exception (error via cancel or read failure) |
| 248 | // Closed is terminal, kj::Exception is implicitly terminal via ErrorState. |
| 249 | using State = StateMachine<TerminalStates<Closed>, |
| 250 | ErrorState<kj::Exception>, |
| 251 | ActiveState<Open>, |
| 252 | Open, |
| 253 | Closed, |
| 254 | kj::Exception>; |
| 255 | State state; |
| 256 | |
| 257 | kj::Rc<WeakRef<ReadableStreamSourceJsAdapter>> selfRef; |
| 258 | }; |
| 259 | |
| 260 | // =============================================================================================== |
| 261 | |
| 262 | // Adapts a ReadableStream to a KJ-friendly interface. |
| 263 | // The adapter fully wraps and consumes the ReadableStream instance, |
| 264 | // using a ReadableStreamDefaultReader to pull data from it. |
| 265 | // When the adapter is destroyed or canceled, the reader is canceled |
| 266 | // and both the reader and the stream references are dropped. Critically, |
| 267 | // the stream is not usable after ownership is transferred to this adapter. |
| 268 | // Initializing the adapter will fail if the stream is already locked or |
| 269 | // disturbed. |
| 270 | // |
| 271 | // If the adapter is dropped, or canceled while there are pending reads, |
| 272 | // the pending reads will be rejected with the same exception as the cancel. |
| 273 | // Because JavaScript promises are not cancelable, reads that are in progress |
| 274 | // won't be aborted immediately but the results will be ignored when they |
| 275 | // complete and a best-effort will be made to interrupt the read as soon as |
| 276 | // possible. If the stream is already closed, reads will complete immediately |
| 277 | // with 0 bytes read. If the stream errors, reads will reject with the same |
| 278 | // exception. |
| 279 | // |
| 280 | // The minRead contract is enforced. The adapter will attempt to read at |
| 281 | // least minBytes on each read, under the isolate lock. If the stream ends |
| 282 | // before minBytes can be satisfied, the read will complete with whatever |
| 283 | // bytes were available and the adapter will remember that the stream is |
| 284 | // closed. |
| 285 | // |
| 286 | // Concurrent/overlapping reads are not allowed. If a read is already |
| 287 | // pending, further read attempts will be rejected. |
| 288 | // |
| 289 | // While the caller is expected to follow the ReadableStreamSource contract |
| 290 | // and keep the adapter and buffer alive until the read promises resolve, |
| 291 | // there are some protections in place to avoid use-after-free if the caller |
| 292 | // drops the adapter. There's nothing we can do if the caller drops the |
| 293 | // buffer, however, so that is still a hard requirement. |
| 294 | // TODO(safety): This can be made safer by having read take a kj::Array |
| 295 | // as input instead of a raw pointer and size, then having the read return |
| 296 | // the filled in Array after the read completes, but that's a larger refactor. |
| 297 | class ReadableSourceKjAdapter final: public ReadableSource { |
| 298 | public: |
| 299 | enum class MinReadPolicy { |
| 300 | // The read will complete as soon as at least minBytes have been read, |
| 301 | // even if more bytes are available and the buffer is not full. This |
| 302 | // may result in more read calls (keeping in mind that each read needs |
| 303 | // to acquire the isolate lock) but may keep the stream flowing more. |
| 304 | IMMEDIATE, |
| 305 | // The read will attempt to fill the entire buffer until either |
| 306 | // maxBytes, the stream ends, or we determine the buffer is "full enough". |
| 307 | // This will result in fewer read calls (and thus grabbing the isolate |
| 308 | // lock less often) but may result in higher latency for each read. |
| 309 | OPPORTUNISTIC, |
| 310 | }; |
| 311 | struct Options { |
| 312 | MinReadPolicy minReadPolicy; |
| 313 | }; |
| 314 | |
| 315 | ReadableSourceKjAdapter(jsg::Lock& js, |
| 316 | IoContext& ioContext, |
| 317 | jsg::Ref<ReadableStream> stream, |
| 318 | Options options = {.minReadPolicy = MinReadPolicy::OPPORTUNISTIC}); |
| 319 | ~ReadableSourceKjAdapter() noexcept(false); |
| 320 | |
| 321 | // Attempts to read at least minBytes and up to maxBytes into the provided |
| 322 | // buffer. The returned promise resolves with the actual number of bytes read, |
| 323 | // which may be less than minBytes if the stream is fully consumed. |
| 324 | // |
| 325 | // If the stream is already closed, the returned promise resolves |
| 326 | // immediately with 0. If the stream is canceled or errors, the returned |
| 327 | // promise rejects with the same exception. |
| 328 | // |
| 329 | // minBytes must be less than or equal to maxBytes and greater than zero. |
| 330 | // If any values outside that range are provided, minBytes will be clamped |
| 331 | // to the range [1, maxBytes]. |
| 332 | // |
| 333 | // Per the contact of tryRead, it is the caller's responsibility to ensure |
| 334 | // that both the buffer and this adapter remain alive until the returned |
| 335 | // promise resolves! It is also the caller's responsibility to ensure that |
| 336 | // buffer is at least maxBytes in length. However, there are some protections |
| 337 | // implemented to avoid use-after-free if the adapter is dropped while a read |
| 338 | // is in progress. |
| 339 | // |
| 340 | // The returned promise will never resolve with more than maxBytes. |
| 341 | kj::Promise<size_t> read(kj::ArrayPtr<kj::byte> buffer, size_t minBytes) override; |
| 342 | |
| 343 | // Reads all remaining bytes from the stream and returns them. |
| 344 | kj::Promise<kj::Array<const kj::byte>> readAllBytes(size_t limit) override; |
| 345 | |
| 346 | // Reads all remaining bytes from the stream and returns them as a string. |
| 347 | kj::Promise<kj::String> readAllText(size_t limit) override; |
| 348 | |
| 349 | // Fully consume the stream and write it to the provided WritableStreamSink. |
| 350 | // If "end" is true, the output stream will be ended once the input |
| 351 | // stream is fully consumed. |
| 352 | // Per the contract of pumpTo, it is the caller's responsibility to ensure |
| 353 | // that both the WritableStreamSink and this adapter remain alive until |
| 354 | // the returned promise resolves! |
| 355 | kj::Promise<DeferredProxy<void>> pumpTo(WritableSink& output, EndAfterPump end) override; |
| 356 | |
| 357 | // If the stream is still active, tries to get the total length, |
| 358 | // if known. If the length is not known, the encoding does not |
| 359 | // match the encoding of the underlying stream, or the stream is closed |
| 360 | // or errored, returns kj::none. |
| 361 | kj::Maybe<size_t> tryGetLength(StreamEncoding encoding) override; |
| 362 | |
| 363 | // Cancels the underlying source if it is still active. |
| 364 | void cancel(kj::Exception reason) override; |
| 365 | |
| 366 | StreamEncoding getEncoding() override { |
| 367 | // Our underlying ReadableStream produces non-encoded bytes (for now) |
| 368 | return StreamEncoding::IDENTITY; |
| 369 | }; |
| 370 | |
| 371 | Tee tee(size_t limit) override; |
| 372 | |
| 373 | struct ReadContext; |
| 374 | KJ_DECLARE_NON_POLYMORPHIC(ReadContext); |
| 375 | |
| 376 | private: |
| 377 | struct Active; |
| 378 | KJ_DECLARE_NON_POLYMORPHIC(Active); |
| 379 | struct KjClosed { |
| 380 | static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj; |
| 381 | }; |
| 382 | struct KjOpen { |
| 383 | static constexpr kj::StringPtr NAME KJ_UNUSED = "open"_kj; |
| 384 | kj::Own<Active> active; |
| 385 | }; |
| 386 | |
| 387 | // State machine for tracking readable source adapter lifecycle: |
| 388 | // KjOpen -> KjClosed (normal close) |
| 389 | // KjOpen -> kj::Exception (error via cancel or read failure) |
| 390 | // KjClosed is terminal, kj::Exception is implicitly terminal via ErrorState. |
| 391 | using KjState = StateMachine<TerminalStates<KjClosed>, |
| 392 | ErrorState<kj::Exception>, |
| 393 | ActiveState<KjOpen>, |
| 394 | KjOpen, |
| 395 | KjClosed, |
| 396 | kj::Exception>; |
| 397 | KjState state; |
| 398 | const Options options; |
| 399 | kj::Rc<WeakRef<ReadableSourceKjAdapter>> selfRef; |
| 400 | |
| 401 | // Checks if the inner Active state is Canceling or Canceled. |
| 402 | // If so, transitions to error state and throws the exception. |
| 403 | void throwIfCancelingOrCanceled(Active& active); |
| 404 | |
| 405 | // Checks if the inner Active state is Canceling or Canceled. |
| 406 | // If so, transitions to error state and returns the exception. |
| 407 | kj::Maybe<kj::Exception> checkCancelingOrCanceled(Active& active); |
| 408 | |
| 409 | kj::Promise<size_t> readImpl(Active& active, kj::ArrayPtr<kj::byte> buffer, size_t minBytes); |
| 410 | |
| 411 | static kj::Promise<void> pumpToImpl( |
| 412 | kj::Own<Active> active, WritableSink& output, EndAfterPump end); |
| 413 | static jsg::Promise<kj::Own<ReadContext>> readInternal( |
| 414 | jsg::Lock& js, kj::Own<ReadContext> context, MinReadPolicy minReadPolicy); |
| 415 | |
| 416 | template <typename T> |
| 417 | kj::Promise<kj::Array<T>> readAllImpl(size_t limit); |
| 418 | |
| 419 | struct CancelationToken final { |
| 420 | kj::Rc<WeakRef<CancelationToken>> selfRef; |
| 421 | inline CancelationToken() |
| 422 | : selfRef(kj::rc<WeakRef<CancelationToken>>(kj::Badge<CancelationToken>{}, *this)) {} |
| 423 | inline ~CancelationToken() { |
| 424 | selfRef->invalidate(); |
| 425 | } |
| 426 | inline kj::Rc<WeakRef<CancelationToken>> getWeakRef() { |
| 427 | return selfRef.addRef(); |
| 428 | } |
| 429 | }; |
| 430 | |
| 431 | template <typename T> |
| 432 | static jsg::Promise<kj::Array<T>> readAllReadImpl(jsg::Lock& js, |
| 433 | IoOwn<Active> active, |
| 434 | kj::Vector<T> accumulated, |
| 435 | size_t limit, |
| 436 | kj::Rc<WeakRef<CancelationToken>> cancelationToken); |
| 437 | }; |
| 438 | |
| 439 | } // namespace workerd::api::streams |