Skip to content
File

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

24.0 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "writable.h"
6 
7#include <workerd/api/system-streams.h>
8#include <workerd/api/worker-rpc.h>
9#include <workerd/io/features.h>
10 
11namespace workerd::api {
12 
13WritableStreamDefaultWriter::WritableStreamDefaultWriter()
14 : ioContext(tryGetIoContext()),
15 state(WriterState::create<Initial>()) {}
16 
17WritableStreamDefaultWriter::~WritableStreamDefaultWriter() noexcept(false) {
18 KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) {
19 attached.stream->getController().releaseWriter(*this, kj::none);
20 }
21}
22 
23jsg::Ref<WritableStreamDefaultWriter> WritableStreamDefaultWriter::constructor(
24 jsg::Lock& js, jsg::Ref<WritableStream> stream) {
25 JSG_REQUIRE(
26 !stream->isLocked(), TypeError, "This WritableStream is currently locked to a writer.");
27 auto writer = js.alloc<WritableStreamDefaultWriter>();
28 writer->lockToStream(js, *stream);
29 return kj::mv(writer);
30}
31 
32jsg::Promise<void> WritableStreamDefaultWriter::abort(
33 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) {
34 assertAttachedOrTerminal();
35 if (state.is<Released>()) {
36 return js.rejectedPromise<void>(
37 js.v8TypeError("This WritableStream writer has been released."_kj));
38 }
39 if (state.is<Closed>()) {
40 return js.resolvedPromise();
41 }
42 auto& attached = state.requireActiveUnsafe();
43 // In some edge cases, this writer is the last thing holding a strong
44 // reference to the stream. Calling abort can cause the writers strong
45 // reference to be cleared, so let's make sure we keep a reference to
46 // the stream at least until the call to abort completes.
47 auto ref = attached.stream.addRef();
48 return attached.stream->getController().abort(js, reason);
49}
50 
51void WritableStreamDefaultWriter::attach(jsg::Lock& js,
52 WritableStreamController& controller,
53 jsg::Promise<void> closedPromise,
54 jsg::Promise<void> readyPromise) {
55 KJ_ASSERT(state.is<Initial>());
56 state.transitionTo<Attached>(controller.addRef());
57 this->closedPromise = kj::mv(closedPromise);
58 replaceReadyPromise(js, kj::mv(readyPromise));
59}
60 
61jsg::Promise<void> WritableStreamDefaultWriter::close(jsg::Lock& js) {
62 assertAttachedOrTerminal();
63 if (state.is<Released>()) {
64 return js.rejectedPromise<void>(
65 js.v8TypeError("This WritableStream writer has been released."_kj));
66 }
67 if (state.is<Closed>()) {
68 return js.rejectedPromise<void>(js.v8TypeError("This WritableStream has been closed."_kj));
69 }
70 auto& attached = state.requireActiveUnsafe();
71 // In some edge cases, this writer is the last thing holding a strong
72 // reference to the stream. Calling close can cause the writers strong
73 // reference to be cleared, so let's make sure we keep a reference to
74 // the stream at least until the call to close completes.
75 auto ref = attached.stream.addRef();
76 return attached.stream->getController().close(js);
77}
78 
79void WritableStreamDefaultWriter::detach() {
80 // Only transition from Attached to Closed.
81 // All other states (Initial, Closed, Released) are no-ops.
82 if (state.isActive()) {
83 state.transitionTo<Closed>();
84 }
85}
86 
87jsg::MemoizedIdentity<jsg::Promise<void>>& WritableStreamDefaultWriter::getClosed() {
88 return KJ_ASSERT_NONNULL(closedPromise, "the writer was never attached to a stream");
89}
90 
91kj::Maybe<int> WritableStreamDefaultWriter::getDesiredSize() {
92 assertAttachedOrTerminal();
93 if (state.is<Released>()) {
94 JSG_FAIL_REQUIRE(TypeError, "This WritableStream writer has been released.");
95 }
96 if (state.is<Closed>()) {
97 return 0;
98 }
99 auto& attached = state.requireActiveUnsafe();
100 return attached.stream->getController().getDesiredSize();
101}
102 
103jsg::MemoizedIdentity<jsg::Promise<void>>& WritableStreamDefaultWriter::getReady() {
104 return KJ_ASSERT_NONNULL(readyPromise, "the writer was never attached to a stream");
105}
106 
107kj::Maybe<jsg::Promise<void>> WritableStreamDefaultWriter::isReady(jsg::Lock& js) {
108 return readyPromisePending.map([&](jsg::Promise<void>& p) { return p.whenResolved(js); });
109}
110 
111void WritableStreamDefaultWriter::lockToStream(jsg::Lock& js, WritableStream& stream) {
112 KJ_ASSERT(!stream.isLocked());
113 KJ_ASSERT(stream.getController().lockWriter(js, *this));
114}
115 
116void WritableStreamDefaultWriter::releaseLock(jsg::Lock& js) {
117 // TODO(soon): Releasing the lock should cancel any pending writes.
118 assertAttachedOrTerminal();
119 // Closed and Released states are no-ops.
120 KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) {
121 // In some edge cases, this writer is the last thing holding a strong
122 // reference to the stream. Calling releaseWriter can cause the writers
123 // strong reference to be cleared, so let's make sure we keep a reference
124 // to the stream at least until the call to releaseLock completes.
125 auto ref = attached.stream.addRef();
126 attached.stream->getController().releaseWriter(*this, js);
127 state.transitionTo<Released>();
128 }
129}
130 
131void WritableStreamDefaultWriter::replaceReadyPromise(
132 jsg::Lock& js, jsg::Promise<void> readyPromise) {
133 this->readyPromisePending = kj::mv(readyPromise);
134 this->readyPromise = KJ_ASSERT_NONNULL(this->readyPromisePending).whenResolved(js);
135}
136 
137jsg::Promise<void> WritableStreamDefaultWriter::write(
138 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> chunk) {
139 assertAttachedOrTerminal();
140 if (state.is<Released>()) {
141 return js.rejectedPromise<void>(
142 js.v8TypeError("This WritableStream writer has been released."_kj));
143 }
144 if (state.is<Closed>()) {
145 return js.rejectedPromise<void>(js.v8TypeError("This WritableStream has been closed."_kj));
146 }
147 auto& attached = state.requireActiveUnsafe();
148 return attached.stream->getController().write(js, chunk);
149}
150 
151jsg::JsString WritableStream::inspectState(jsg::Lock& js) {
152 if (controller->isErrored()) {
153 return js.strIntern("errored");
154 } else if (controller->isErroring(js) != kj::none) {
155 return js.strIntern("erroring");
156 } else if (controller->isClosedOrClosing()) {
157 return js.strIntern("closed");
158 } else {
159 return js.strIntern("writable");
160 }
161}
162 
163bool WritableStream::inspectExpectsBytes() {
164 return controller->isByteOriented();
165}
166 
167void WritableStreamDefaultWriter::visitForGc(jsg::GcVisitor& visitor) {
168 KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) {
169 visitor.visit(attached.stream);
170 }
171 visitor.visit(closedPromise, readyPromise);
172}
173 
174// ======================================================================================
175 
176WritableStream::WritableStream(IoContext& ioContext,
177 kj::Own<WritableStreamSink> sink,
178 kj::Maybe<kj::Own<ByteStreamObserver>> maybeObserver,
179 kj::Maybe<uint64_t> maybeHighWaterMark,
180 kj::Maybe<jsg::Promise<void>> maybeClosureWaitable)
181 : WritableStream(newWritableStreamInternalController(ioContext,
182 kj::mv(sink),
183 kj::mv(maybeObserver),
184 maybeHighWaterMark,
185 kj::mv(maybeClosureWaitable))) {}
186 
187WritableStream::WritableStream(kj::Own<WritableStreamController> controller)
188 : ioContext(tryGetIoContext()),
189 controller(kj::mv(controller)) {
190 getController().setOwnerRef(*this);
191}
192 
193jsg::Ref<WritableStream> WritableStream::addRef() {
194 return JSG_THIS;
195}
196 
197void WritableStream::visitForGc(jsg::GcVisitor& visitor) {
198 visitor.visit(getController());
199}
200 
201bool WritableStream::isLocked() {
202 return getController().isLockedToWriter();
203}
204 
205WritableStreamController& WritableStream::getController() {
206 return *controller;
207}
208 
209kj::Own<WritableStreamSink> WritableStream::removeSink(jsg::Lock& js) {
210 return JSG_REQUIRE_NONNULL(getController().removeSink(js), TypeError,
211 "This WritableStream does not have a WritableStreamSink");
212}
213 
214void WritableStream::detach(jsg::Lock& js) {
215 getController().detach(js);
216}
217 
218jsg::Promise<void> WritableStream::abort(
219 jsg::Lock& js, jsg::Optional<v8::Local<v8::Value>> reason) {
220 if (isLocked()) {
221 return js.rejectedPromise<void>(
222 js.v8TypeError("This WritableStream is currently locked to a writer."_kj));
223 }
224 return getController().abort(js, reason);
225}
226 
227jsg::Promise<void> WritableStream::close(jsg::Lock& js) {
228 if (isLocked()) {
229 return js.rejectedPromise<void>(
230 js.v8TypeError("This WritableStream is currently locked to a writer."_kj));
231 }
232 return getController().close(js);
233}
234 
235jsg::Promise<void> WritableStream::flush(jsg::Lock& js) {
236 if (isLocked()) {
237 return js.rejectedPromise<void>(
238 js.v8TypeError("This WritableStream is currently locked to a writer."_kj));
239 }
240 return getController().flush(js);
241}
242 
243jsg::Ref<WritableStreamDefaultWriter> WritableStream::getWriter(jsg::Lock& js) {
244 return WritableStreamDefaultWriter::constructor(js, JSG_THIS);
245}
246 
247jsg::Ref<WritableStream> WritableStream::constructor(jsg::Lock& js,
248 jsg::Optional<UnderlyingSink> underlyingSink,
249 jsg::Optional<StreamQueuingStrategy> queuingStrategy) {
250 JSG_REQUIRE(FeatureFlags::get(js).getStreamsJavaScriptControllers(), Error,
251 "To use the new WritableStream() constructor, enable the "
252 "streams_enable_constructors compatibility flag. "
253 "Refer to the docs for more information: https://developers.cloudflare.com/workers/platform/compatibility-dates/#compatibility-flags");
254 auto controller = newWritableStreamJsController();
255 // We account for the memory usage of the WritableStream and its controller together because their
256 // lifetimes are identical and memory accounting itself has a memory overhead.
257 auto stream = js.allocAccounted<WritableStream>(
258 sizeof(WritableStream) + controller->jsgGetMemorySelfSize(), kj::mv(controller));
259 stream->getController().setup(js, kj::mv(underlyingSink), kj::mv(queuingStrategy));
260 return kj::mv(stream);
261}
262 
263namespace {
264 
265// Wrapper around `WritableStreamSink` that makes it suitable for passing off to capnp RPC.
266class WritableStreamRpcAdapter final: public capnp::ExplicitEndOutputStream {
267 public:
268 WritableStreamRpcAdapter(kj::Own<WritableStreamSink> inner): inner(kj::mv(inner)) {}
269 ~WritableStreamRpcAdapter() noexcept(false) {
270 weakRef->invalidate();
271 doneFulfiller->fulfill();
272 }
273 
274 // Returns a promise that resolves when the stream is dropped. If the promise is canceled before
275 // that, the stream is revoked.
276 kj::Promise<void> waitForCompletionOrRevoke() {
277 auto paf = kj::newPromiseAndFulfiller<void>();
278 doneFulfiller = kj::mv(paf.fulfiller);
279 
280 return paf.promise.attach(kj::defer([weakRef = weakRef->addRef()]() mutable {
281 KJ_IF_SOME(obj, weakRef->tryGet()) {
282 // Stream is still alive, revoke it.
283 if (!obj.canceler.isEmpty()) {
284 obj.canceler.cancel(cancellationException());
285 }
286 obj.inner = kj::none;
287 }
288 }));
289 }
290 
291 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
292 return canceler.wrap(getInner().write(buffer));
293 }
294 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
295 return canceler.wrap(getInner().write(pieces));
296 }
297 
298 // TODO(perf): We can't properly implement tryPumpFrom(), which means that Cap'n Proto will
299 // be unable to perform path shortening if the underlying stream turns out to be another capnp
300 // stream. This isn't a huge deal, but might be nice to enable someday. It may require
301 // significant refactoring of streams.
302 
303 kj::Promise<void> whenWriteDisconnected() override {
304 // TODO(someday): WritableStreamSink doesn't give us a way to implement this.
305 return kj::NEVER_DONE;
306 }
307 
308 kj::Promise<void> end() override {
309 return canceler.wrap(getInner().end());
310 }
311 
312 private:
313 kj::Maybe<kj::Own<WritableStreamSink>> inner;
314 kj::Canceler canceler;
315 kj::Own<kj::PromiseFulfiller<void>> doneFulfiller;
316 kj::Own<WeakRef<WritableStreamRpcAdapter>> weakRef =
317 kj::refcounted<WeakRef<WritableStreamRpcAdapter>>(
318 kj::Badge<WritableStreamRpcAdapter>(), *this);
319 
320 WritableStreamSink& getInner() {
321 return *KJ_UNWRAP_OR(inner, { kj::throwFatalException(cancellationException()); });
322 }
323 
324 static kj::Exception cancellationException() {
325 return JSG_KJ_EXCEPTION(DISCONNECTED, Error,
326 "WritableStream received over RPC was disconnected because the remote execution context "
327 "has endeded.");
328 }
329};
330 
331// In order to support JavaScript-backed WritableStreams that do not have a backing
332// WritableStreamSink, we need an alternative version of the WritableStreamRpcAdapter
333// that will arrange to acquire the isolate lock when necessary to perform writes
334// directly on the WritableStreamController. Note that this approach is necessarily
335// a lot slower
336class WritableStreamJsRpcAdapter final: public capnp::ExplicitEndOutputStream {
337 public:
338 WritableStreamJsRpcAdapter(IoContext& context, jsg::Ref<WritableStreamDefaultWriter> writer)
339 : context(context),
340 writer(kj::mv(writer)) {}
341 
342 ~WritableStreamJsRpcAdapter() noexcept(false) {
343 weakRef->invalidate();
344 doneFulfiller->fulfill();
345 
346 // If the stream was not explicitly ended and the writer still exists at this point,
347 // then we should trigger calling the abort algorithm on the stream. Sadly, there's a
348 // bit of an incompatibility with kj::AsyncOutputStream and the standard definition of
349 // WritableStream in that AsyncOutputStream has no specific way to explicitly signal that
350 // the stream is being aborted due to a particular reason.
351 //
352 // On the remote side, because it is using a WritableStreamSink implementation, when that
353 // side is aborted, all it does is record the reason and drop the stream. It does not
354 // propagate the reason back to this side. So, we have to do the best we can here. Our
355 // assumption is that once the stream is dropped, if it has not been explicitly ended and
356 // the writer still exists, then the writer should be aborted. This is not perfect because
357 // we cannot propagate the actual reason why it was aborted.
358 //
359 // Note also that there is no guarantee that the abort will actually run if the context
360 // is being torn down. Some WritableStream implementations might use the abort algorithm
361 // to clean things up or perform logging in the case of an error. Care needs to be taken
362 // in this situation or the user code might end up with bugs. Need to see if there's a
363 // better solution.
364 //
365 // TODO(someday): If the remote end can be updated to propagate the abort, then we can
366 // hopefully improve the situation here.
367 if (!ended) {
368 KJ_IF_SOME(writer, this->writer) {
369 context.addTask(context.run([writer = kj::mv(writer), exception = cancellationException()](
370 Worker::Lock& lock) mutable {
371 jsg::Lock& js = lock;
372 auto ex = js.exceptionToJs(kj::mv(exception));
373 return IoContext::current().awaitJs(lock, writer->abort(lock, ex.getHandle(js)));
374 }));
375 }
376 }
377 }
378 
379 // Returns a promise that resolves when the stream is dropped. If the promise is canceled before
380 // that, the stream is revoked.
381 kj::Promise<void> waitForCompletionOrRevoke() {
382 auto paf = kj::newPromiseAndFulfiller<void>();
383 doneFulfiller = kj::mv(paf.fulfiller);
384 
385 return paf.promise.attach(kj::defer([weakRef = weakRef->addRef()]() mutable {
386 KJ_IF_SOME(obj, weakRef->tryGet()) {
387 // Stream is still alive, revoke it.
388 if (!obj.canceler.isEmpty()) {
389 obj.canceler.cancel(cancellationException());
390 }
391 auto w = kj::mv(obj.writer);
392 KJ_IF_SOME(writer, w) {
393 obj.context.addTask(
394 obj.context.run([writer = kj::mv(writer), exception = cancellationException()](
395 Worker::Lock& lock) mutable {
396 jsg::Lock& js = lock;
397 auto ex = js.exceptionToJs(kj::mv(exception));
398 return IoContext::current().awaitJs(lock, writer->abort(lock, ex.getHandle(js)));
399 }));
400 }
401 }
402 }));
403 }
404 
405 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
406 if (writer == kj::none) {
407 return KJ_EXCEPTION(FAILED, "Write after stream has been closed.");
408 }
409 if (buffer == nullptr) return kj::READY_NOW;
410 return canceler.wrap(context.run([this, buffer](Worker::Lock& lock) mutable {
411 auto& writer = getInner();
412 auto source = KJ_ASSERT_NONNULL(jsg::BufferSource::tryAlloc(lock, buffer.size()));
413 source.asArrayPtr().copyFrom(buffer);
414 return context.awaitJs(lock, writer.write(lock, source.getHandle(lock)));
415 }));
416 }
417 
418 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
419 if (writer == kj::none) {
420 return KJ_EXCEPTION(FAILED, "Write after stream has been closed.");
421 }
422 auto amount = 0;
423 for (auto& piece: pieces) {
424 amount += piece.size();
425 }
426 if (amount == 0) return kj::READY_NOW;
427 return canceler.wrap(context.run([this, amount, pieces](Worker::Lock& lock) mutable {
428 auto& writer = getInner();
429 // Sadly, we have to allocate and copy here. Our received set of buffers are only
430 // guaranteed to live until the returned promise is resolved, but the application code
431 // may hold onto the ArrayBuffer for longer. We need to make sure that the backing store
432 // for the ArrayBuffer remains valid.
433 auto source = KJ_ASSERT_NONNULL(jsg::BufferSource::tryAlloc(lock, amount));
434 auto ptr = source.asArrayPtr();
435 for (auto& piece: pieces) {
436 KJ_DASSERT(ptr.size() > 0);
437 KJ_DASSERT(piece.size() <= ptr.size());
438 if (piece.size() == 0) continue;
439 ptr.first(piece.size()).copyFrom(piece);
440 ptr = ptr.slice(piece.size());
441 }
442 
443 return context.awaitJs(lock, writer.write(lock, source.getHandle(lock)));
444 }));
445 }
446 
447 // TODO(perf): We can't properly implement tryPumpFrom(), which means that Cap'n Proto will
448 // be unable to perform path shortening if the underlying stream turns out to be another capnp
449 // stream. This isn't a huge deal, but might be nice to enable someday. It may require
450 // significant refactoring of streams.
451 
452 kj::Promise<void> whenWriteDisconnected() override {
453 // TODO(soon): We might be able to support this by following the writer.closed promise,
454 // which becomes resolved when the writer is used to close the stream, or rejects when
455 // the stream has errored. However, currently, we don't have an easy way to do this.
456 //
457 // The Writer's getClosed() method returns a jsg::MemoizedIdentity<jsg::Promise<void>>.
458 // jsg::MemoizedIdentity lazily converts the jsg::Promise into a v8::Promise once it
459 // passes through the type wrapper. It does not give us any way to consistently get
460 // at the underlying jsg::Promise<void> or the mapped v8::Promise. We would need to
461 // capture a TypeHandler in here and convert each time to one or the other, then
462 // attach our continuation. It's doable but a bit of a pain.
463 //
464 // For now, let's handle this the same as WritableStreamRpcAdapter and just return a
465 // never done.
466 return kj::NEVER_DONE;
467 }
468 
469 kj::Promise<void> end() override {
470 if (writer == kj::none) {
471 return KJ_EXCEPTION(FAILED, "End after stream has been closed.");
472 }
473 ended = true;
474 return canceler.wrap(context.run([this](Worker::Lock& lock) mutable {
475 return context.awaitJs(lock, getInner().close(lock));
476 }));
477 }
478 
479 private:
480 IoContext& context;
481 kj::Maybe<jsg::Ref<WritableStreamDefaultWriter>> writer;
482 kj::Canceler canceler;
483 kj::Own<kj::PromiseFulfiller<void>> doneFulfiller;
484 kj::Own<WeakRef<WritableStreamJsRpcAdapter>> weakRef =
485 kj::refcounted<WeakRef<WritableStreamJsRpcAdapter>>(
486 kj::Badge<WritableStreamJsRpcAdapter>(), *this);
487 bool ended = false;
488 
489 WritableStreamDefaultWriter& getInner() {
490 KJ_IF_SOME(inner, writer) {
491 return *inner;
492 }
493 kj::throwFatalException(cancellationException());
494 }
495 
496 static kj::Exception cancellationException() {
497 return JSG_KJ_EXCEPTION(DISCONNECTED, Error,
498 "WritableStream received over RPC was disconnected because the remote execution context "
499 "has endeded.");
500 }
501};
502 
503} // namespace
504 
505void WritableStream::serialize(jsg::Lock& js, jsg::Serializer& serializer) {
506 // Serialize by effectively creating a `JsRpcStub` around this object and serializing that.
507 // Except we don't actually want to do _exactly_ that, because we do not want to actually create
508 // a `JsRpcStub` locally. So do the important parts of `JsRpcStub::constructor()` followed by
509 // `JsRpcStub::serialize()`.
510 
511 auto& handler = JSG_REQUIRE_NONNULL(serializer.getExternalHandler(), DOMDataCloneError,
512 "WritableStream can only be serialized for RPC.");
513 auto externalHandler = dynamic_cast<RpcSerializerExternalHandler*>(&handler);
514 JSG_REQUIRE(externalHandler != nullptr, DOMDataCloneError,
515 "WritableStream can only be serialized for RPC.");
516 
517 IoContext& ioctx = IoContext::current();
518 
519 // TODO(soon): Support JS-backed WritableStreams. Currently this only supports native streams
520 // and IdentityTransformStream, since only they are backed by WritableStreamSink.
521 
522 KJ_IF_SOME(sink, getController().removeSink(js)) {
523 // NOTE: We're counting on `removeSink()`, to check that the stream is not locked and other
524 // common checks. It's important we don't modify the WritableStream before this call.
525 auto encoding = sink->disownEncodingResponsibility();
526 auto wrapper = kj::heap<WritableStreamRpcAdapter>(kj::mv(sink));
527 
528 // Make sure this stream will be revoked if the IoContext ends.
529 ioctx.addTask(wrapper->waitForCompletionOrRevoke().attach(ioctx.registerPendingEvent()));
530 
531 auto capnpStream = ioctx.getByteStreamFactory().kjToCapnp(kj::mv(wrapper));
532 
533 externalHandler->write([capnpStream = kj::mv(capnpStream), encoding](
534 rpc::JsValue::External::Builder builder) mutable {
535 auto ws = builder.initWritableStream();
536 ws.setByteStream(kj::mv(capnpStream));
537 ws.setEncoding(encoding);
538 });
539 } else {
540 // TODO(soon): Support disownEncodingResponsibility with JS-backed streams
541 
542 // NOTE: We're counting on `getWriter()` to check that the stream is not locked and other
543 // common checks. It's important we don't modify the WritableStream before this call.
544 auto wrapper = kj::heap<WritableStreamJsRpcAdapter>(ioctx, getWriter(js));
545 
546 // Make sure this stream will be revoked if the IoContext ends.
547 ioctx.addTask(wrapper->waitForCompletionOrRevoke().attach(ioctx.registerPendingEvent()));
548 
549 auto capnpStream = ioctx.getByteStreamFactory().kjToCapnp(kj::mv(wrapper));
550 
551 externalHandler->write(
552 [capnpStream = kj::mv(capnpStream)](rpc::JsValue::External::Builder builder) mutable {
553 auto ws = builder.initWritableStream();
554 ws.setByteStream(kj::mv(capnpStream));
555 ws.setEncoding(StreamEncoding::IDENTITY);
556 });
557 }
558}
559 
560jsg::Ref<WritableStream> WritableStream::deserialize(
561 jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer) {
562 auto& handler = KJ_REQUIRE_NONNULL(
563 deserializer.getExternalHandler(), "got WritableStream on non-RPC serialized object?");
564 auto externalHandler = dynamic_cast<RpcDeserializerExternalHandler*>(&handler);
565 KJ_REQUIRE(externalHandler != nullptr, "got WritableStream on non-RPC serialized object?");
566 
567 auto reader = externalHandler->read();
568 KJ_REQUIRE(reader.isWritableStream(), "external table slot type doesn't match serialization tag");
569 
570 auto ws = reader.getWritableStream();
571 auto encoding = ws.getEncoding();
572 
573 KJ_REQUIRE(
574 static_cast<uint>(encoding) < capnp::Schema::from<StreamEncoding>().getEnumerants().size(),
575 "unknown StreamEncoding received from peer");
576 
577 IoContext& ioctx = IoContext::current();
578 auto stream = ioctx.getByteStreamFactory().capnpToKjExplicitEnd(ws.getByteStream());
579 auto sink = newSystemStream(kj::mv(stream), encoding, ioctx);
580 
581 return js.alloc<WritableStream>(
582 ioctx, kj::mv(sink), ioctx.getMetrics().tryCreateWritableByteStreamObserver());
583}
584 
585void WritableStreamDefaultWriter::visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
586 KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) {
587 tracker.trackField("attached", attached.stream);
588 }
589 tracker.trackField("closedPromise", closedPromise);
590 tracker.trackField("readyPromise", readyPromise);
591}
592 
593void WritableStream::visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
594 tracker.trackField("controller", controller);
595}
596 
597} // namespace workerd::api