Skip to content
File

Blob: src/workerd/api/streams/writable-sink-adapter.c++

25.8 KB
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 
8namespace 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.
18struct 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 
130WritableStreamSinkJsAdapter::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 
141WritableStreamSinkJsAdapter::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 
152WritableStreamSinkJsAdapter::~WritableStreamSinkJsAdapter() noexcept(false) {
153 selfRef->invalidate();
154}
155 
156kj::Maybe<const kj::Exception&> WritableStreamSinkJsAdapter::isErrored() {
157 return state.tryGetErrorUnsafe();
158}
159 
160bool WritableStreamSinkJsAdapter::isClosed() {
161 return state.is<Closed>();
162}
163 
164bool WritableStreamSinkJsAdapter::isClosing() {
165 return state.whenActiveOr([](Open& open) { return open.active->closePending; }, false);
166}
167 
168kj::Maybe<ssize_t> WritableStreamSinkJsAdapter::getDesiredSize() {
169 KJ_IF_SOME(open, state.tryGetActiveUnsafe()) {
170 return open.active->getDesiredSize();
171 }
172 return kj::none;
173}
174 
175jsg::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 
305jsg::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.
342jsg::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.
389void 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 
401void WritableStreamSinkJsAdapter::abort(jsg::Lock& js, const jsg::JsValue& reason) {
402 abort(js.exceptionToKj(reason));
403}
404 
405void 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 
416void 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 
424bool WritableStreamSinkJsAdapter::BackpressureState::isWaiting() const {
425 return readyResolver != kj::none;
426}
427 
428jsg::Promise<void> WritableStreamSinkJsAdapter::BackpressureState::getReady(jsg::Lock& js) {
429 return ready.whenResolved(js);
430}
431 
432jsg::MemoizedIdentity<jsg::Promise<void>>& WritableStreamSinkJsAdapter::BackpressureState::
433 getReadyStable() {
434 return readyWatcher;
435}
436 
437WritableStreamSinkJsAdapter::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 
445void 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 
457void 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 
467WritableStreamSinkJsAdapter::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 
475jsg::Promise<void> WritableStreamSinkJsAdapter::getReady(jsg::Lock& js) {
476 return backpressureState.getReady(js);
477}
478 
479jsg::MemoizedIdentity<jsg::Promise<void>>& WritableStreamSinkJsAdapter::getReadyStable() {
480 return backpressureState.getReadyStable();
481}
482 
483void WritableStreamSinkJsAdapter::visitForGc(jsg::GcVisitor& visitor) {
484 visitor.visit(
485 backpressureState.readyResolver, backpressureState.ready, backpressureState.readyWatcher);
486}
487 
488void 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 
494kj::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 
504struct 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 
527namespace {
528jsg::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 
534WritableStreamSinkKjAdapter::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 
540WritableStreamSinkKjAdapter::Active::~Active() noexcept(false) {
541 abort(KJ_EXCEPTION(DISCONNECTED, "WritableStreamSinkKjAdapter is canceled."));
542}
543 
544void 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 
557WritableStreamSinkKjAdapter::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 
563WritableStreamSinkKjAdapter::~WritableStreamSinkKjAdapter() noexcept(false) {
564 selfRef->invalidate();
565}
566 
567kj::Promise<void> WritableStreamSinkKjAdapter::write(kj::ArrayPtr<const byte> buffer) {
568 auto pieces = kj::arr(buffer);
569 co_await write(pieces);
570}
571 
572kj::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 
642kj::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 
679void 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