File
Blob: src/workerd/api/streams/writable-sink-adapter.c++
| 1 | #include "writable-sink-adapter.h" |
| 2 | |
| 3 | #include "writable.h" |
| 4 | |
| 5 | #include <workerd/api/system-streams.h> |
| 6 | #include <workerd/util/checked-queue.h> |
| 7 | |
| 8 | namespace workerd::api::streams { |
| 9 | |
| 10 | // The Active state maintains a queue of tasks, such as write or flush operations. Each task |
| 11 | // contains a promise-returning function object and a fulfiller. When the first task is |
| 12 | // enqueued, the active state begins processing the queue asynchronously. Each function |
| 13 | // is invoked in order, its promise awaited, and the result passed to the fulfiller. The |
| 14 | // fulfiller notifies the code which enqueued the task that the task has completed. In |
| 15 | // this way, read and close operations are safely executed in serial, even if one operation |
| 16 | // is called before the previous completes. This mechanism satisfies KJ's restriction on |
| 17 | // concurrent operations on streams. |
| 18 | struct WritableStreamSinkJsAdapter::Active final { |
| 19 | struct Task { |
| 20 | kj::Function<kj::Promise<void>()> task; |
| 21 | kj::Own<kj::PromiseFulfiller<void>> fulfiller; |
| 22 | kj::Maybe<kj::Promise<void>> maybeOutputLock; |
| 23 | |
| 24 | Task(kj::Function<kj::Promise<void>()> task, |
| 25 | kj::Own<kj::PromiseFulfiller<void>> fulfiller, |
| 26 | kj::Maybe<kj::Promise<void>> maybeOutputLock = kj::none) |
| 27 | : task(kj::mv(task)), |
| 28 | fulfiller(kj::mv(fulfiller)), |
| 29 | maybeOutputLock(kj::mv(maybeOutputLock)) {} |
| 30 | KJ_DISALLOW_COPY_AND_MOVE(Task); |
| 31 | }; |
| 32 | using TaskQueue = workerd::util::Queue<kj::Own<Task>>; |
| 33 | |
| 34 | kj::Own<WritableSink> sink; |
| 35 | const Options options; |
| 36 | kj::Canceler canceler; |
| 37 | TaskQueue queue; |
| 38 | bool aborted = false; |
| 39 | bool running = false; |
| 40 | bool closePending = false; |
| 41 | size_t bytesInFlight = 0; |
| 42 | kj::Maybe<kj::Exception> pendingAbort; |
| 43 | |
| 44 | Active(kj::Own<WritableSink> sink, Options options) |
| 45 | : sink(kj::mv(sink)), |
| 46 | options(kj::mv(options)) { |
| 47 | KJ_DASSERT(this->sink.get() != nullptr, "WritableStreamSink cannot be null"); |
| 48 | } |
| 49 | |
| 50 | KJ_DISALLOW_COPY_AND_MOVE(Active); |
| 51 | ~Active() noexcept(false) { |
| 52 | // When the Active is dropped, we cancel any remaining pending writes and |
| 53 | // abort the sink. |
| 54 | abort(KJ_EXCEPTION(FAILED, "jsg.Error: Writable stream is canceled or closed.")); |
| 55 | |
| 56 | // Check invariants for safety. |
| 57 | // 1. Our canceler should be empty because we canceled it. |
| 58 | KJ_DASSERT(canceler.isEmpty()); |
| 59 | // 2. The write queue should be empty. |
| 60 | KJ_DASSERT(queue.empty()); |
| 61 | } |
| 62 | |
| 63 | // Explicitly cancel all in-flight and pending tasks in the queue. |
| 64 | // This is a non-op if cancel has already been called. |
| 65 | void abort(kj::Exception&& exception) { |
| 66 | if (aborted) return; |
| 67 | aborted = true; |
| 68 | // 1. Cancel our in-flight "runLoop", if any. |
| 69 | pendingAbort = exception.clone(); |
| 70 | canceler.cancel(exception.clone()); |
| 71 | // 2. Drop our queue of pending tasks. |
| 72 | queue.drainTo( |
| 73 | [&exception](kj::Own<Task>&& task) { task->fulfiller->reject(exception.clone()); }); |
| 74 | // 3. Abort and drop the sink itself. We're done with it. |
| 75 | sink->abort(kj::mv(exception)); |
| 76 | auto dropped KJ_UNUSED = kj::mv(sink); |
| 77 | } |
| 78 | |
| 79 | // Get the desired size based on the configured high water mark and |
| 80 | // the number of bytes currently in flight. |
| 81 | ssize_t getDesiredSize() const { |
| 82 | return options.highWaterMark - bytesInFlight; |
| 83 | } |
| 84 | |
| 85 | kj::Promise<void> enqueue(kj::Function<kj::Promise<void>()> task) { |
| 86 | KJ_DASSERT(!aborted, "cannot enqueue tasks on an aborted queue"); |
| 87 | auto paf = kj::newPromiseAndFulfiller<void>(); |
| 88 | auto& ioContext = IoContext::current(); |
| 89 | queue.push(kj::heap<Task>( |
| 90 | kj::mv(task), kj::mv(paf.fulfiller), ioContext.waitForOutputLocksIfNecessary())); |
| 91 | if (!running) { |
| 92 | ioContext.addTask(canceler.wrap(run())); |
| 93 | } |
| 94 | return kj::mv(paf.promise); |
| 95 | } |
| 96 | |
| 97 | kj::Promise<void> run() { |
| 98 | KJ_DEFER(running = false); |
| 99 | running = true; |
| 100 | while (!queue.empty() && !aborted) { |
| 101 | auto task = KJ_ASSERT_NONNULL(queue.pop()); |
| 102 | KJ_DEFER({ |
| 103 | if (task->fulfiller->isWaiting()) { |
| 104 | KJ_IF_SOME(pending, pendingAbort) { |
| 105 | task->fulfiller->reject(kj::mv(pending)); |
| 106 | } else { |
| 107 | task->fulfiller->reject(KJ_EXCEPTION(DISCONNECTED, "Task was canceled.")); |
| 108 | } |
| 109 | } |
| 110 | }); |
| 111 | bool taskFailed = false; |
| 112 | try { |
| 113 | KJ_IF_SOME(lock, task->maybeOutputLock) { |
| 114 | co_await lock; |
| 115 | } |
| 116 | co_await task->task(); |
| 117 | task->fulfiller->fulfill(); |
| 118 | } catch (...) { |
| 119 | auto ex = kj::getCaughtExceptionAsKj(); |
| 120 | task->fulfiller->reject(kj::mv(ex)); |
| 121 | taskFailed = true; |
| 122 | } |
| 123 | // If the task failed, we exit the loop. We're going to abort the |
| 124 | // entire remaining queue anyway so there's no point in continuing. |
| 125 | if (taskFailed) co_return; |
| 126 | } |
| 127 | } |
| 128 | }; |
| 129 | |
| 130 | WritableStreamSinkJsAdapter::WritableStreamSinkJsAdapter( |
| 131 | jsg::Lock& js, IoContext& ioContext, kj::Own<WritableSink> sink, kj::Maybe<Options> options) |
| 132 | : state(State::create<Open>( |
| 133 | ioContext.addObject(kj::heap<Active>(kj::mv(sink), kj::mv(options).orDefault({}))))), |
| 134 | backpressureState(newBackpressureState(js)), |
| 135 | selfRef(kj::rc<WeakRef<WritableStreamSinkJsAdapter>>( |
| 136 | kj::Badge<WritableStreamSinkJsAdapter>{}, *this)) { |
| 137 | // We want the initial backpressure state to be "ready". |
| 138 | backpressureState.release(js); |
| 139 | } |
| 140 | |
| 141 | WritableStreamSinkJsAdapter::WritableStreamSinkJsAdapter(jsg::Lock& js, |
| 142 | IoContext& ioContext, |
| 143 | kj::Own<kj::AsyncOutputStream> stream, |
| 144 | StreamEncoding encoding, |
| 145 | kj::Maybe<Options> options) |
| 146 | : WritableStreamSinkJsAdapter(js, |
| 147 | ioContext, |
| 148 | newIoContextWrappedWritableSink( |
| 149 | ioContext, newEncodedWritableSink(encoding, kj::mv(stream))), |
| 150 | kj::mv(options)) {} |
| 151 | |
| 152 | WritableStreamSinkJsAdapter::~WritableStreamSinkJsAdapter() noexcept(false) { |
| 153 | selfRef->invalidate(); |
| 154 | } |
| 155 | |
| 156 | kj::Maybe<const kj::Exception&> WritableStreamSinkJsAdapter::isErrored() { |
| 157 | return state.tryGetErrorUnsafe(); |
| 158 | } |
| 159 | |
| 160 | bool WritableStreamSinkJsAdapter::isClosed() { |
| 161 | return state.is<Closed>(); |
| 162 | } |
| 163 | |
| 164 | bool WritableStreamSinkJsAdapter::isClosing() { |
| 165 | return state.whenActiveOr([](Open& open) { return open.active->closePending; }, false); |
| 166 | } |
| 167 | |
| 168 | kj::Maybe<ssize_t> WritableStreamSinkJsAdapter::getDesiredSize() { |
| 169 | KJ_IF_SOME(open, state.tryGetActiveUnsafe()) { |
| 170 | return open.active->getDesiredSize(); |
| 171 | } |
| 172 | return kj::none; |
| 173 | } |
| 174 | |
| 175 | jsg::Promise<void> WritableStreamSinkJsAdapter::write(jsg::Lock& js, const jsg::JsValue& value) { |
| 176 | KJ_IF_SOME(exc, state.tryGetErrorUnsafe()) { |
| 177 | // Really should not have been called if errored but just in case, |
| 178 | // return a rejected promise. |
| 179 | return js.rejectedPromise<void>(js.exceptionToJs(exc.clone())); |
| 180 | } |
| 181 | |
| 182 | if (state.is<Closed>()) { |
| 183 | // Really should not have been called if closed but just in case, |
| 184 | // return a rejected promise. |
| 185 | return js.rejectedPromise<void>(js.typeError("Write after close is not allowed")); |
| 186 | } |
| 187 | |
| 188 | auto& open = state.requireActiveUnsafe(); |
| 189 | // Dereference the IoOwn once to get the active state. |
| 190 | auto& active = *open.active; |
| 191 | |
| 192 | // If close is pending, we cannot accept any more writes. |
| 193 | if (active.closePending) { |
| 194 | auto exc = js.typeError("Write after close is not allowed"); |
| 195 | return js.rejectedPromise<void>(exc); |
| 196 | } |
| 197 | |
| 198 | // Ok, we are in a writable state, there are no pending closes. |
| 199 | // Let's process our data and write it! |
| 200 | auto& ioContext = IoContext::current(); |
| 201 | |
| 202 | // We know that a WritableStreamSink only accepts bytes, so we need to |
| 203 | // verify that the value is a source of bytes. We accept three possible |
| 204 | // types: ArrayBuffer, ArrayBufferView, and String. If it is a string, |
| 205 | // we convert it to UTF-8 bytes. Anything else is an error. |
| 206 | if (value.isArrayBufferView() || value.isArrayBuffer() || value.isSharedArrayBuffer()) { |
| 207 | // We can just wrap the value with a jsg::BufferSource and write it. |
| 208 | jsg::BufferSource source(js, value); |
| 209 | if (active.options.detachOnWrite && source.canDetach(js)) { |
| 210 | // Detach from the original ArrayBuffer... |
| 211 | // ... and re-wrap it with a new BufferSource that we own. |
| 212 | source = jsg::BufferSource(js, source.detach(js)); |
| 213 | } |
| 214 | |
| 215 | // Zero-length writes are a no-op. |
| 216 | if (source.size() == 0) { |
| 217 | return js.resolvedPromise(); |
| 218 | } |
| 219 | |
| 220 | active.bytesInFlight += source.size(); |
| 221 | maybeSignalBackpressure(js); |
| 222 | // Enqueue the actual write operation into the write queue. We pass in |
| 223 | // two lambdas, one that does the actual write, and one that handles |
| 224 | // errors. If the write fails, we need to transition the adapter to the |
| 225 | // errored state. If the write succeeds, we need to decrement the |
| 226 | // bytesInFlight counter. |
| 227 | // |
| 228 | // The promise returned by enqueue is not the actual write promise but |
| 229 | // a branch forked off of it. We wrap that with a JS promise that waits |
| 230 | // for it to complete. Once it does, we check if we can release backpressure. |
| 231 | // This has to be done within an Isolate lock because we need to be able |
| 232 | // to resolve or reject the JS promises. If the write fails, we instead |
| 233 | // abort the backpressure state. |
| 234 | // |
| 235 | // This slight indirection does mean that the backpressure state change |
| 236 | // may be slightly delayed after the actual write completes but that's |
| 237 | // ok. |
| 238 | // |
| 239 | // Capturing active by reference here is safe because the lambda is |
| 240 | // held by the write queue, which is itself held by Active. If active |
| 241 | // is destroyed, the write queue is destroyed along with the lambda. |
| 242 | auto promise = |
| 243 | active.enqueue(kj::coCapture([&active, source = kj::mv(source)]() -> kj::Promise<void> { |
| 244 | co_await active.sink->write(source.asArrayPtr()); |
| 245 | active.bytesInFlight -= source.size(); |
| 246 | })); |
| 247 | return ioContext |
| 248 | .awaitIo(js, kj::mv(promise), [self = selfRef.addRef()](jsg::Lock& js) { |
| 249 | // Why do we need a weak ref here? Well, because this is a JavaScript |
| 250 | // promise continuation. It is possible that the kj::Own holding our |
| 251 | // adapter can be dropped while we are waiting for the continuation |
| 252 | // to run. If that happens, we don't want to delay cleanup of the |
| 253 | // adapter just because of backpressure state management that would |
| 254 | // not be needed anymore, so we use a weak ref to update the backpressure |
| 255 | // state only if we are still alive. |
| 256 | self->runIfAlive( |
| 257 | [&](WritableStreamSinkJsAdapter& self) { self.maybeReleaseBackpressure(js); }); |
| 258 | }).catch_(js, [self = selfRef.addRef()](jsg::Lock& js, jsg::Value exception) { |
| 259 | auto error = jsg::JsValue(exception.getHandle(js)); |
| 260 | self->runIfAlive([&](WritableStreamSinkJsAdapter& self) { |
| 261 | self.abort(js, error); |
| 262 | self.backpressureState.abort(js, error); |
| 263 | }); |
| 264 | js.throwException(kj::mv(exception)); |
| 265 | }); |
| 266 | } else if (value.isString()) { |
| 267 | // Also super easy! Let's just convert the string to UTF-8 |
| 268 | auto str = value.toString(js); |
| 269 | |
| 270 | // Zero-length writes are a no-op. |
| 271 | if (str.size() == 0) { |
| 272 | return js.resolvedPromise(); |
| 273 | } |
| 274 | |
| 275 | active.bytesInFlight += str.size(); |
| 276 | // Make sure to account for the memory used by the string while the |
| 277 | // write is in-flight/pending |
| 278 | auto accounting = js.getExternalMemoryAdjustment(str.size()); |
| 279 | maybeSignalBackpressure(js); |
| 280 | // Just like above, enqueue the write operation into the write queue, |
| 281 | // ensuring that we handle both the success and failure cases. |
| 282 | auto promise = active.enqueue(kj::coCapture( |
| 283 | [&active, str = kj::mv(str), accounting = kj::mv(accounting)]() -> kj::Promise<void> { |
| 284 | co_await active.sink->write(str.asBytes()); |
| 285 | active.bytesInFlight -= str.size(); |
| 286 | })); |
| 287 | return ioContext |
| 288 | .awaitIo(js, kj::mv(promise), [self = selfRef.addRef()](jsg::Lock& js) { |
| 289 | self->runIfAlive( |
| 290 | [&](WritableStreamSinkJsAdapter& self) { self.maybeReleaseBackpressure(js); }); |
| 291 | }).catch_(js, [self = selfRef.addRef()](jsg::Lock& js, jsg::Value exception) { |
| 292 | auto error = jsg::JsValue(exception.getHandle(js)); |
| 293 | self->runIfAlive([&](WritableStreamSinkJsAdapter& self) { |
| 294 | self.abort(js, error); |
| 295 | self.backpressureState.abort(js, error); |
| 296 | }); |
| 297 | js.throwException(kj::mv(exception)); |
| 298 | }); |
| 299 | } |
| 300 | |
| 301 | auto err = js.typeError("This WritableStream only supports writing byte types."_kj); |
| 302 | return js.rejectedPromise<void>(err); |
| 303 | } |
| 304 | |
| 305 | jsg::Promise<void> WritableStreamSinkJsAdapter::flush(jsg::Lock& js) { |
| 306 | KJ_IF_SOME(exc, state.tryGetErrorUnsafe()) { |
| 307 | // Really should not have been called if errored but just in case, |
| 308 | // return a rejected promise. |
| 309 | return js.rejectedPromise<void>(js.exceptionToJs(exc.clone())); |
| 310 | } |
| 311 | |
| 312 | if (state.is<Closed>()) { |
| 313 | // Really should not have been called if closed but just in case, |
| 314 | // return a rejected promise. |
| 315 | return js.rejectedPromise<void>(js.typeError("Flush after close is not allowed")); |
| 316 | } |
| 317 | |
| 318 | auto& open = state.requireActiveUnsafe(); |
| 319 | // Dereference the IoOwn once to get the active state. |
| 320 | auto& active = *open.active; |
| 321 | |
| 322 | // If close is pending, we cannot accept any more writes. |
| 323 | if (active.closePending) { |
| 324 | auto exc = js.typeError("Flush after close is not allowed"); |
| 325 | return js.rejectedPromise<void>(exc); |
| 326 | } |
| 327 | |
| 328 | // Ok, we are in a writable state, there are no pending closes. |
| 329 | // Let's enqueue our flush signal. |
| 330 | auto& ioContext = IoContext::current(); |
| 331 | // Flushing is really just a non-op write. We enqueue a no-op task |
| 332 | // into the write queue and wait for it to complete. |
| 333 | auto promise = active.enqueue([]() -> kj::Promise<void> { |
| 334 | // Non-op. |
| 335 | return kj::READY_NOW; |
| 336 | }); |
| 337 | return ioContext.awaitIo(js, kj::mv(promise)); |
| 338 | } |
| 339 | |
| 340 | // Transitions the adapter into the closing state. Once the write queue |
| 341 | // is empty, we will close the sink and transition to the closed state. |
| 342 | jsg::Promise<void> WritableStreamSinkJsAdapter::end(jsg::Lock& js) { |
| 343 | KJ_IF_SOME(exc, state.tryGetErrorUnsafe()) { |
| 344 | // Really should not have been called if errored but just in case, |
| 345 | // return a rejected promise. |
| 346 | return js.rejectedPromise<void>(js.exceptionToJs(exc.clone())); |
| 347 | } |
| 348 | |
| 349 | if (state.is<Closed>()) { |
| 350 | // We are already in a closed state. This is a no-op. This really |
| 351 | // should not have been called if closed but just in case, return |
| 352 | // a resolved promise. |
| 353 | return js.resolvedPromise(); |
| 354 | } |
| 355 | |
| 356 | auto& open = state.requireActiveUnsafe(); |
| 357 | auto& ioContext = IoContext::current(); |
| 358 | auto& active = *open.active; |
| 359 | |
| 360 | if (active.closePending) { |
| 361 | return js.rejectedPromise<void>(js.typeError("Close already pending, cannot close again.")); |
| 362 | } |
| 363 | |
| 364 | active.closePending = true; |
| 365 | auto promise = active.enqueue( |
| 366 | kj::coCapture([&active]() -> kj::Promise<void> { co_await active.sink->end(); })); |
| 367 | |
| 368 | return ioContext |
| 369 | .awaitIo(js, kj::mv(promise), [self = selfRef.addRef()](jsg::Lock& js) { |
| 370 | // While nothing at this point should be actually waiting on the ready promise, |
| 371 | // we should still resolve it just in case. |
| 372 | self->runIfAlive([&](WritableStreamSinkJsAdapter& self) { |
| 373 | self.state.transitionTo<Closed>(); |
| 374 | self.maybeReleaseBackpressure(js); |
| 375 | }); |
| 376 | }).catch_(js, [self = selfRef.addRef()](jsg::Lock& js, jsg::Value&& exception) { |
| 377 | // Likewise, while nothing should be waiting on the ready promise, we |
| 378 | // should still reject it just in case. |
| 379 | auto error = jsg::JsValue(exception.getHandle(js)); |
| 380 | self->runIfAlive([&](WritableStreamSinkJsAdapter& self) { |
| 381 | self.abort(js, error); |
| 382 | self.backpressureState.abort(js, error); |
| 383 | }); |
| 384 | js.throwException(kj::mv(exception)); |
| 385 | }); |
| 386 | } |
| 387 | |
| 388 | // Transitions the adapter to the errored state, even if we are already closed. |
| 389 | void WritableStreamSinkJsAdapter::abort(kj::Exception&& exception) { |
| 390 | // If we are in an active state, we need to cancel any in-flight and pending |
| 391 | // operations in the active write queue *before* we transition to the errored |
| 392 | // state. This ensures that any pending writes are interrupted and do not |
| 393 | // complete. |
| 394 | KJ_IF_SOME(open, state.tryGetActiveUnsafe()) { |
| 395 | open.active->abort(exception.clone()); |
| 396 | } |
| 397 | // Use forceTransitionTo because abort can be called from any state. |
| 398 | state.forceTransitionTo<kj::Exception>(kj::mv(exception)); |
| 399 | } |
| 400 | |
| 401 | void WritableStreamSinkJsAdapter::abort(jsg::Lock& js, const jsg::JsValue& reason) { |
| 402 | abort(js.exceptionToKj(reason)); |
| 403 | } |
| 404 | |
| 405 | void WritableStreamSinkJsAdapter::BackpressureState::abort( |
| 406 | jsg::Lock& js, const jsg::JsValue& reason) { |
| 407 | // Backpressure signaling is being aborted, likely because the adapter |
| 408 | // transitioned to the errored state. Reject the ready promise with |
| 409 | // the given reason. |
| 410 | KJ_IF_SOME(resolver, readyResolver) { |
| 411 | resolver.reject(js, reason); |
| 412 | readyResolver = kj::none; |
| 413 | } |
| 414 | } |
| 415 | |
| 416 | void WritableStreamSinkJsAdapter::BackpressureState::release(jsg::Lock& js) { |
| 417 | // The backppressure has been released. Resolve the ready promise. |
| 418 | KJ_IF_SOME(resolver, readyResolver) { |
| 419 | resolver.resolve(js); |
| 420 | readyResolver = kj::none; |
| 421 | } |
| 422 | } |
| 423 | |
| 424 | bool WritableStreamSinkJsAdapter::BackpressureState::isWaiting() const { |
| 425 | return readyResolver != kj::none; |
| 426 | } |
| 427 | |
| 428 | jsg::Promise<void> WritableStreamSinkJsAdapter::BackpressureState::getReady(jsg::Lock& js) { |
| 429 | return ready.whenResolved(js); |
| 430 | } |
| 431 | |
| 432 | jsg::MemoizedIdentity<jsg::Promise<void>>& WritableStreamSinkJsAdapter::BackpressureState:: |
| 433 | getReadyStable() { |
| 434 | return readyWatcher; |
| 435 | } |
| 436 | |
| 437 | WritableStreamSinkJsAdapter::BackpressureState::BackpressureState( |
| 438 | jsg::Promise<void>::Resolver&& resolver, |
| 439 | jsg::Promise<void>&& promise, |
| 440 | jsg::MemoizedIdentity<jsg::Promise<void>>&& watcher) |
| 441 | : readyResolver(kj::mv(resolver)), |
| 442 | ready(kj::mv(promise)), |
| 443 | readyWatcher(kj::mv(watcher)) {} |
| 444 | |
| 445 | void WritableStreamSinkJsAdapter::maybeSignalBackpressure(jsg::Lock& js) { |
| 446 | // We should only be signaling backpressure if we are in an active state. |
| 447 | state.requireActiveUnsafe(); |
| 448 | // Indicate that backpressure is being applied. If we are already in a |
| 449 | // backpressure state (isWaiting() is true), this is a no-op. |
| 450 | if (!backpressureState.isWaiting()) { |
| 451 | // We signal backpressure by replacing the backpressure state. |
| 452 | // This replaces the JS promises and resolvers with a new set. |
| 453 | backpressureState = newBackpressureState(js); |
| 454 | } |
| 455 | } |
| 456 | |
| 457 | void WritableStreamSinkJsAdapter::maybeReleaseBackpressure(jsg::Lock& js) { |
| 458 | KJ_IF_SOME(open, state.tryGetActiveUnsafe()) { |
| 459 | if (open.active->getDesiredSize() > 0) { |
| 460 | // The desired size is now > 0, so we can release backpressure. |
| 461 | // If backpressure is already released or aborted, this is a non-op. |
| 462 | backpressureState.release(js); |
| 463 | } |
| 464 | } |
| 465 | } |
| 466 | |
| 467 | WritableStreamSinkJsAdapter::BackpressureState WritableStreamSinkJsAdapter::newBackpressureState( |
| 468 | jsg::Lock& js) { |
| 469 | jsg::PromiseResolverPair<void> pair = js.newPromiseAndResolver<void>(); |
| 470 | pair.promise.markAsHandled(js); |
| 471 | auto watcher = jsg::MemoizedIdentity<jsg::Promise<void>>(pair.promise.whenResolved(js)); |
| 472 | return BackpressureState(kj::mv(pair.resolver), kj::mv(pair.promise), kj::mv(watcher)); |
| 473 | } |
| 474 | |
| 475 | jsg::Promise<void> WritableStreamSinkJsAdapter::getReady(jsg::Lock& js) { |
| 476 | return backpressureState.getReady(js); |
| 477 | } |
| 478 | |
| 479 | jsg::MemoizedIdentity<jsg::Promise<void>>& WritableStreamSinkJsAdapter::getReadyStable() { |
| 480 | return backpressureState.getReadyStable(); |
| 481 | } |
| 482 | |
| 483 | void WritableStreamSinkJsAdapter::visitForGc(jsg::GcVisitor& visitor) { |
| 484 | visitor.visit( |
| 485 | backpressureState.readyResolver, backpressureState.ready, backpressureState.readyWatcher); |
| 486 | } |
| 487 | |
| 488 | void WritableStreamSinkJsAdapter::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 489 | tracker.trackField("backpressureState.readyResolver", backpressureState.readyResolver); |
| 490 | tracker.trackField("backpressureState.ready", backpressureState.ready); |
| 491 | tracker.trackField("backpressureState.readyWatcher", backpressureState.readyWatcher); |
| 492 | } |
| 493 | |
| 494 | kj::Maybe<const WritableStreamSinkJsAdapter::Options&> WritableStreamSinkJsAdapter::getOptions() { |
| 495 | KJ_IF_SOME(open, state.tryGetActiveUnsafe()) { |
| 496 | return open.active->options; |
| 497 | } else { |
| 498 | return kj::none; |
| 499 | } |
| 500 | } |
| 501 | |
| 502 | // ================================================================================================ |
| 503 | |
| 504 | struct WritableStreamSinkKjAdapter::Active { |
| 505 | IoContext& ioContext; |
| 506 | jsg::Ref<WritableStream> stream; |
| 507 | jsg::Ref<WritableStreamDefaultWriter> writer; |
| 508 | kj::Canceler canceler; |
| 509 | |
| 510 | // The contract of WritableStreamSink is that there can only be one |
| 511 | // write in-flight at a time. |
| 512 | bool writePending = false; |
| 513 | |
| 514 | bool closePending = false; |
| 515 | kj::Maybe<kj::Exception> pendingAbort; |
| 516 | |
| 517 | // Prevent abort() from being called multiple times. |
| 518 | bool aborted = false; |
| 519 | |
| 520 | Active(jsg::Lock& js, IoContext& ioContext, jsg::Ref<WritableStream> stream); |
| 521 | KJ_DISALLOW_COPY_AND_MOVE(Active); |
| 522 | ~Active() noexcept(false); |
| 523 | |
| 524 | void abort(kj::Exception reason); |
| 525 | }; |
| 526 | |
| 527 | namespace { |
| 528 | jsg::Ref<WritableStreamDefaultWriter> initWriter(jsg::Lock& js, jsg::Ref<WritableStream>& stream) { |
| 529 | JSG_REQUIRE(!stream->isLocked(), TypeError, "WritableStream is locked."); |
| 530 | return stream->getWriter(js); |
| 531 | } |
| 532 | } // namespace |
| 533 | |
| 534 | WritableStreamSinkKjAdapter::Active::Active( |
| 535 | jsg::Lock& js, IoContext& ioContext, jsg::Ref<WritableStream> stream) |
| 536 | : ioContext(ioContext), |
| 537 | stream(kj::mv(stream)), |
| 538 | writer(initWriter(js, this->stream)) {} |
| 539 | |
| 540 | WritableStreamSinkKjAdapter::Active::~Active() noexcept(false) { |
| 541 | abort(KJ_EXCEPTION(DISCONNECTED, "WritableStreamSinkKjAdapter is canceled.")); |
| 542 | } |
| 543 | |
| 544 | void WritableStreamSinkKjAdapter::Active::abort(kj::Exception reason) { |
| 545 | if (aborted) return; |
| 546 | aborted = true; |
| 547 | canceler.cancel(reason.clone()); |
| 548 | ioContext.addTask(ioContext.run([writable = kj::mv(stream), writer = kj::mv(writer), |
| 549 | exception = reason.clone()](jsg::Lock& js) mutable { |
| 550 | auto& ioContext = IoContext::current(); |
| 551 | auto error = js.exceptionToJsValue(kj::mv(exception)); |
| 552 | auto promise = writer->abort(js, error.getHandle(js)); |
| 553 | return ioContext.awaitJs(js, kj::mv(promise)); |
| 554 | })); |
| 555 | } |
| 556 | |
| 557 | WritableStreamSinkKjAdapter::WritableStreamSinkKjAdapter( |
| 558 | jsg::Lock& js, IoContext& ioContext, jsg::Ref<WritableStream> stream) |
| 559 | : state(KjState::create<KjOpen>(kj::heap<Active>(js, ioContext, kj::mv(stream)))), |
| 560 | selfRef(kj::rc<WeakRef<WritableStreamSinkKjAdapter>>( |
| 561 | kj::Badge<WritableStreamSinkKjAdapter>{}, *this)) {} |
| 562 | |
| 563 | WritableStreamSinkKjAdapter::~WritableStreamSinkKjAdapter() noexcept(false) { |
| 564 | selfRef->invalidate(); |
| 565 | } |
| 566 | |
| 567 | kj::Promise<void> WritableStreamSinkKjAdapter::write(kj::ArrayPtr<const byte> buffer) { |
| 568 | auto pieces = kj::arr(buffer); |
| 569 | co_await write(pieces); |
| 570 | } |
| 571 | |
| 572 | kj::Promise<void> WritableStreamSinkKjAdapter::write( |
| 573 | kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) { |
| 574 | KJ_IF_SOME(exc, state.tryGetErrorUnsafe()) { |
| 575 | kj::throwFatalException(exc.clone()); |
| 576 | } |
| 577 | |
| 578 | if (state.is<KjClosed>()) { |
| 579 | KJ_FAIL_REQUIRE("Cannot write after close."); |
| 580 | } |
| 581 | |
| 582 | auto& open = state.requireActiveUnsafe(); |
| 583 | auto& active = *open.active; |
| 584 | KJ_REQUIRE(!active.writePending, "Cannot have multiple concurrent writes."); |
| 585 | KJ_IF_SOME(exception, active.pendingAbort) { |
| 586 | auto exc = exception.clone(); |
| 587 | state.forceTransitionTo<kj::Exception>(exc.clone()); |
| 588 | return kj::mv(exc); |
| 589 | } |
| 590 | if (active.closePending) { |
| 591 | state.transitionTo<KjClosed>(); |
| 592 | KJ_FAIL_REQUIRE("Cannot write after close."); |
| 593 | } |
| 594 | active.writePending = true; |
| 595 | |
| 596 | return active.canceler |
| 597 | .wrap(active.ioContext.run([self = selfRef.addRef(), writer = active.writer.addRef(), |
| 598 | pieces = pieces](jsg::Lock& js) mutable -> kj::Promise<void> { |
| 599 | size_t totalAmount = 0; |
| 600 | for (auto piece: pieces) { |
| 601 | totalAmount += piece.size(); |
| 602 | } |
| 603 | if (totalAmount == 0) { |
| 604 | return kj::READY_NOW; |
| 605 | } |
| 606 | |
| 607 | // We collapse our pieces into a single ArrayBuffer for efficiency. The |
| 608 | // WritableStream API has no concept of a vector write, so each write |
| 609 | // would incur the overhead of a separate promise and microtask checkpoint. |
| 610 | // By collapsing into a single write we reduce that overhead. |
| 611 | auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, totalAmount); |
| 612 | auto ptr = backing.asArrayPtr(); |
| 613 | for (auto piece: pieces) { |
| 614 | ptr.first(piece.size()).copyFrom(piece); |
| 615 | ptr = ptr.slice(piece.size()); |
| 616 | } |
| 617 | jsg::BufferSource source(js, kj::mv(backing)); |
| 618 | |
| 619 | auto ready = KJ_ASSERT_NONNULL(writer->isReady(js)); |
| 620 | auto promise = |
| 621 | ready.then(js, [writer = writer.addRef(), source = kj::mv(source)](jsg::Lock& js) mutable { |
| 622 | return writer->write(js, source.getHandle(js)); |
| 623 | }); |
| 624 | return IoContext::current().awaitJs(js, kj::mv(promise)); |
| 625 | })).then([self = selfRef.addRef()]() { |
| 626 | self->runIfAlive([&](WritableStreamSinkKjAdapter& self) { |
| 627 | KJ_IF_SOME(open, self.state.tryGetActiveUnsafe()) { |
| 628 | open.active->writePending = false; |
| 629 | } |
| 630 | }); |
| 631 | }, [self = selfRef.addRef()](kj::Exception exception) { |
| 632 | self->runIfAlive([&](WritableStreamSinkKjAdapter& self) { |
| 633 | KJ_IF_SOME(open, self.state.tryGetActiveUnsafe()) { |
| 634 | open.active->writePending = false; |
| 635 | open.active->pendingAbort = exception.clone(); |
| 636 | } |
| 637 | }); |
| 638 | kj::throwFatalException(kj::mv(exception)); |
| 639 | }); |
| 640 | } |
| 641 | |
| 642 | kj::Promise<void> WritableStreamSinkKjAdapter::end() { |
| 643 | KJ_IF_SOME(exc, state.tryGetErrorUnsafe()) { |
| 644 | return exc.clone(); |
| 645 | } |
| 646 | |
| 647 | if (state.is<KjClosed>()) { |
| 648 | return kj::READY_NOW; |
| 649 | } |
| 650 | |
| 651 | auto& open = state.requireActiveUnsafe(); |
| 652 | auto& active = *open.active; |
| 653 | KJ_REQUIRE(!active.writePending, "Cannot have multiple concurrent writes."); |
| 654 | KJ_IF_SOME(exception, active.pendingAbort) { |
| 655 | auto exc = kj::mv(exception); |
| 656 | state.forceTransitionTo<kj::Exception>(exc.clone()); |
| 657 | return kj::mv(exc); |
| 658 | } |
| 659 | if (active.closePending) { |
| 660 | state.transitionTo<KjClosed>(); |
| 661 | return kj::READY_NOW; |
| 662 | } |
| 663 | active.closePending = true; |
| 664 | return active.canceler |
| 665 | .wrap(active.ioContext.run( |
| 666 | [self = selfRef.addRef(), writer = active.writer.addRef()](jsg::Lock& js) mutable { |
| 667 | auto promise = writer->close(js); |
| 668 | return IoContext::current().awaitJs(js, kj::mv(promise)); |
| 669 | })).catch_([self = selfRef.addRef()](kj::Exception exception) { |
| 670 | self->runIfAlive([&](WritableStreamSinkKjAdapter& self) { |
| 671 | KJ_IF_SOME(open, self.state.tryGetActiveUnsafe()) { |
| 672 | open.active->pendingAbort = exception.clone(); |
| 673 | } |
| 674 | }); |
| 675 | kj::throwFatalException(kj::mv(exception)); |
| 676 | }); |
| 677 | } |
| 678 | |
| 679 | void WritableStreamSinkKjAdapter::abort(kj::Exception reason) { |
| 680 | KJ_IF_SOME(open, state.tryGetActiveUnsafe()) { |
| 681 | open.active->abort(reason.clone()); |
| 682 | } |
| 683 | // Use forceTransitionTo because abort can be called from any state. |
| 684 | state.forceTransitionTo<kj::Exception>(kj::mv(reason)); |
| 685 | } |
| 686 | |
| 687 | } // namespace workerd::api::streams |