// Copyright (c) 2017-2022 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #include "writable.h" #include #include #include namespace workerd::api { WritableStreamDefaultWriter::WritableStreamDefaultWriter() : ioContext(tryGetIoContext()), state(WriterState::create()) {} WritableStreamDefaultWriter::~WritableStreamDefaultWriter() noexcept(false) { KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) { attached.stream->getController().releaseWriter(*this, kj::none); } } jsg::Ref WritableStreamDefaultWriter::constructor( jsg::Lock& js, jsg::Ref stream) { JSG_REQUIRE( !stream->isLocked(), TypeError, "This WritableStream is currently locked to a writer."); auto writer = js.alloc(); writer->lockToStream(js, *stream); return kj::mv(writer); } jsg::Promise WritableStreamDefaultWriter::abort( jsg::Lock& js, jsg::Optional> reason) { assertAttachedOrTerminal(); if (state.is()) { return js.rejectedPromise( js.v8TypeError("This WritableStream writer has been released."_kj)); } if (state.is()) { return js.resolvedPromise(); } auto& attached = state.requireActiveUnsafe(); // In some edge cases, this writer is the last thing holding a strong // reference to the stream. Calling abort can cause the writers strong // reference to be cleared, so let's make sure we keep a reference to // the stream at least until the call to abort completes. auto ref = attached.stream.addRef(); return attached.stream->getController().abort(js, reason); } void WritableStreamDefaultWriter::attach(jsg::Lock& js, WritableStreamController& controller, jsg::Promise closedPromise, jsg::Promise readyPromise) { KJ_ASSERT(state.is()); state.transitionTo(controller.addRef()); this->closedPromise = kj::mv(closedPromise); replaceReadyPromise(js, kj::mv(readyPromise)); } jsg::Promise WritableStreamDefaultWriter::close(jsg::Lock& js) { assertAttachedOrTerminal(); if (state.is()) { return js.rejectedPromise( js.v8TypeError("This WritableStream writer has been released."_kj)); } if (state.is()) { return js.rejectedPromise(js.v8TypeError("This WritableStream has been closed."_kj)); } auto& attached = state.requireActiveUnsafe(); // In some edge cases, this writer is the last thing holding a strong // reference to the stream. Calling close can cause the writers strong // reference to be cleared, so let's make sure we keep a reference to // the stream at least until the call to close completes. auto ref = attached.stream.addRef(); return attached.stream->getController().close(js); } void WritableStreamDefaultWriter::detach() { // Only transition from Attached to Closed. // All other states (Initial, Closed, Released) are no-ops. if (state.isActive()) { state.transitionTo(); } } jsg::MemoizedIdentity>& WritableStreamDefaultWriter::getClosed() { return KJ_ASSERT_NONNULL(closedPromise, "the writer was never attached to a stream"); } kj::Maybe WritableStreamDefaultWriter::getDesiredSize() { assertAttachedOrTerminal(); if (state.is()) { JSG_FAIL_REQUIRE(TypeError, "This WritableStream writer has been released."); } if (state.is()) { return 0; } auto& attached = state.requireActiveUnsafe(); return attached.stream->getController().getDesiredSize(); } jsg::MemoizedIdentity>& WritableStreamDefaultWriter::getReady() { return KJ_ASSERT_NONNULL(readyPromise, "the writer was never attached to a stream"); } kj::Maybe> WritableStreamDefaultWriter::isReady(jsg::Lock& js) { return readyPromisePending.map([&](jsg::Promise& p) { return p.whenResolved(js); }); } void WritableStreamDefaultWriter::lockToStream(jsg::Lock& js, WritableStream& stream) { KJ_ASSERT(!stream.isLocked()); KJ_ASSERT(stream.getController().lockWriter(js, *this)); } void WritableStreamDefaultWriter::releaseLock(jsg::Lock& js) { // TODO(soon): Releasing the lock should cancel any pending writes. assertAttachedOrTerminal(); // Closed and Released states are no-ops. KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) { // In some edge cases, this writer is the last thing holding a strong // reference to the stream. Calling releaseWriter can cause the writers // strong reference to be cleared, so let's make sure we keep a reference // to the stream at least until the call to releaseLock completes. auto ref = attached.stream.addRef(); attached.stream->getController().releaseWriter(*this, js); state.transitionTo(); } } void WritableStreamDefaultWriter::replaceReadyPromise( jsg::Lock& js, jsg::Promise readyPromise) { this->readyPromisePending = kj::mv(readyPromise); this->readyPromise = KJ_ASSERT_NONNULL(this->readyPromisePending).whenResolved(js); } jsg::Promise WritableStreamDefaultWriter::write( jsg::Lock& js, jsg::Optional> chunk) { assertAttachedOrTerminal(); if (state.is()) { return js.rejectedPromise( js.v8TypeError("This WritableStream writer has been released."_kj)); } if (state.is()) { return js.rejectedPromise(js.v8TypeError("This WritableStream has been closed."_kj)); } auto& attached = state.requireActiveUnsafe(); return attached.stream->getController().write(js, chunk); } jsg::JsString WritableStream::inspectState(jsg::Lock& js) { if (controller->isErrored()) { return js.strIntern("errored"); } else if (controller->isErroring(js) != kj::none) { return js.strIntern("erroring"); } else if (controller->isClosedOrClosing()) { return js.strIntern("closed"); } else { return js.strIntern("writable"); } } bool WritableStream::inspectExpectsBytes() { return controller->isByteOriented(); } void WritableStreamDefaultWriter::visitForGc(jsg::GcVisitor& visitor) { KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) { visitor.visit(attached.stream); } visitor.visit(closedPromise, readyPromise); } // ====================================================================================== WritableStream::WritableStream(IoContext& ioContext, kj::Own sink, kj::Maybe> maybeObserver, kj::Maybe maybeHighWaterMark, kj::Maybe> maybeClosureWaitable) : WritableStream(newWritableStreamInternalController(ioContext, kj::mv(sink), kj::mv(maybeObserver), maybeHighWaterMark, kj::mv(maybeClosureWaitable))) {} WritableStream::WritableStream(kj::Own controller) : ioContext(tryGetIoContext()), controller(kj::mv(controller)) { getController().setOwnerRef(*this); } jsg::Ref WritableStream::addRef() { return JSG_THIS; } void WritableStream::visitForGc(jsg::GcVisitor& visitor) { visitor.visit(getController()); } bool WritableStream::isLocked() { return getController().isLockedToWriter(); } WritableStreamController& WritableStream::getController() { return *controller; } kj::Own WritableStream::removeSink(jsg::Lock& js) { return JSG_REQUIRE_NONNULL(getController().removeSink(js), TypeError, "This WritableStream does not have a WritableStreamSink"); } void WritableStream::detach(jsg::Lock& js) { getController().detach(js); } jsg::Promise WritableStream::abort( jsg::Lock& js, jsg::Optional> reason) { if (isLocked()) { return js.rejectedPromise( js.v8TypeError("This WritableStream is currently locked to a writer."_kj)); } return getController().abort(js, reason); } jsg::Promise WritableStream::close(jsg::Lock& js) { if (isLocked()) { return js.rejectedPromise( js.v8TypeError("This WritableStream is currently locked to a writer."_kj)); } return getController().close(js); } jsg::Promise WritableStream::flush(jsg::Lock& js) { if (isLocked()) { return js.rejectedPromise( js.v8TypeError("This WritableStream is currently locked to a writer."_kj)); } return getController().flush(js); } jsg::Ref WritableStream::getWriter(jsg::Lock& js) { return WritableStreamDefaultWriter::constructor(js, JSG_THIS); } jsg::Ref WritableStream::constructor(jsg::Lock& js, jsg::Optional underlyingSink, jsg::Optional queuingStrategy) { JSG_REQUIRE(FeatureFlags::get(js).getStreamsJavaScriptControllers(), Error, "To use the new WritableStream() constructor, enable the " "streams_enable_constructors compatibility flag. " "Refer to the docs for more information: https://developers.cloudflare.com/workers/platform/compatibility-dates/#compatibility-flags"); auto controller = newWritableStreamJsController(); // We account for the memory usage of the WritableStream and its controller together because their // lifetimes are identical and memory accounting itself has a memory overhead. auto stream = js.allocAccounted( sizeof(WritableStream) + controller->jsgGetMemorySelfSize(), kj::mv(controller)); stream->getController().setup(js, kj::mv(underlyingSink), kj::mv(queuingStrategy)); return kj::mv(stream); } namespace { // Wrapper around `WritableStreamSink` that makes it suitable for passing off to capnp RPC. class WritableStreamRpcAdapter final: public capnp::ExplicitEndOutputStream { public: WritableStreamRpcAdapter(kj::Own inner): inner(kj::mv(inner)) {} ~WritableStreamRpcAdapter() noexcept(false) { weakRef->invalidate(); doneFulfiller->fulfill(); } // Returns a promise that resolves when the stream is dropped. If the promise is canceled before // that, the stream is revoked. kj::Promise waitForCompletionOrRevoke() { auto paf = kj::newPromiseAndFulfiller(); doneFulfiller = kj::mv(paf.fulfiller); return paf.promise.attach(kj::defer([weakRef = weakRef->addRef()]() mutable { KJ_IF_SOME(obj, weakRef->tryGet()) { // Stream is still alive, revoke it. if (!obj.canceler.isEmpty()) { obj.canceler.cancel(cancellationException()); } obj.inner = kj::none; } })); } kj::Promise write(kj::ArrayPtr buffer) override { return canceler.wrap(getInner().write(buffer)); } kj::Promise write(kj::ArrayPtr> pieces) override { return canceler.wrap(getInner().write(pieces)); } // TODO(perf): We can't properly implement tryPumpFrom(), which means that Cap'n Proto will // be unable to perform path shortening if the underlying stream turns out to be another capnp // stream. This isn't a huge deal, but might be nice to enable someday. It may require // significant refactoring of streams. kj::Promise whenWriteDisconnected() override { // TODO(someday): WritableStreamSink doesn't give us a way to implement this. return kj::NEVER_DONE; } kj::Promise end() override { return canceler.wrap(getInner().end()); } private: kj::Maybe> inner; kj::Canceler canceler; kj::Own> doneFulfiller; kj::Own> weakRef = kj::refcounted>( kj::Badge(), *this); WritableStreamSink& getInner() { return *KJ_UNWRAP_OR(inner, { kj::throwFatalException(cancellationException()); }); } static kj::Exception cancellationException() { return JSG_KJ_EXCEPTION(DISCONNECTED, Error, "WritableStream received over RPC was disconnected because the remote execution context " "has endeded."); } }; // In order to support JavaScript-backed WritableStreams that do not have a backing // WritableStreamSink, we need an alternative version of the WritableStreamRpcAdapter // that will arrange to acquire the isolate lock when necessary to perform writes // directly on the WritableStreamController. Note that this approach is necessarily // a lot slower class WritableStreamJsRpcAdapter final: public capnp::ExplicitEndOutputStream { public: WritableStreamJsRpcAdapter(IoContext& context, jsg::Ref writer) : context(context), writer(kj::mv(writer)) {} ~WritableStreamJsRpcAdapter() noexcept(false) { weakRef->invalidate(); doneFulfiller->fulfill(); // If the stream was not explicitly ended and the writer still exists at this point, // then we should trigger calling the abort algorithm on the stream. Sadly, there's a // bit of an incompatibility with kj::AsyncOutputStream and the standard definition of // WritableStream in that AsyncOutputStream has no specific way to explicitly signal that // the stream is being aborted due to a particular reason. // // On the remote side, because it is using a WritableStreamSink implementation, when that // side is aborted, all it does is record the reason and drop the stream. It does not // propagate the reason back to this side. So, we have to do the best we can here. Our // assumption is that once the stream is dropped, if it has not been explicitly ended and // the writer still exists, then the writer should be aborted. This is not perfect because // we cannot propagate the actual reason why it was aborted. // // Note also that there is no guarantee that the abort will actually run if the context // is being torn down. Some WritableStream implementations might use the abort algorithm // to clean things up or perform logging in the case of an error. Care needs to be taken // in this situation or the user code might end up with bugs. Need to see if there's a // better solution. // // TODO(someday): If the remote end can be updated to propagate the abort, then we can // hopefully improve the situation here. if (!ended) { KJ_IF_SOME(writer, this->writer) { context.addTask(context.run([writer = kj::mv(writer), exception = cancellationException()]( Worker::Lock& lock) mutable { jsg::Lock& js = lock; auto ex = js.exceptionToJs(kj::mv(exception)); return IoContext::current().awaitJs(lock, writer->abort(lock, ex.getHandle(js))); })); } } } // Returns a promise that resolves when the stream is dropped. If the promise is canceled before // that, the stream is revoked. kj::Promise waitForCompletionOrRevoke() { auto paf = kj::newPromiseAndFulfiller(); doneFulfiller = kj::mv(paf.fulfiller); return paf.promise.attach(kj::defer([weakRef = weakRef->addRef()]() mutable { KJ_IF_SOME(obj, weakRef->tryGet()) { // Stream is still alive, revoke it. if (!obj.canceler.isEmpty()) { obj.canceler.cancel(cancellationException()); } auto w = kj::mv(obj.writer); KJ_IF_SOME(writer, w) { obj.context.addTask( obj.context.run([writer = kj::mv(writer), exception = cancellationException()]( Worker::Lock& lock) mutable { jsg::Lock& js = lock; auto ex = js.exceptionToJs(kj::mv(exception)); return IoContext::current().awaitJs(lock, writer->abort(lock, ex.getHandle(js))); })); } } })); } kj::Promise write(kj::ArrayPtr buffer) override { if (writer == kj::none) { return KJ_EXCEPTION(FAILED, "Write after stream has been closed."); } if (buffer == nullptr) return kj::READY_NOW; return canceler.wrap(context.run([this, buffer](Worker::Lock& lock) mutable { auto& writer = getInner(); auto source = KJ_ASSERT_NONNULL(jsg::BufferSource::tryAlloc(lock, buffer.size())); source.asArrayPtr().copyFrom(buffer); return context.awaitJs(lock, writer.write(lock, source.getHandle(lock))); })); } kj::Promise write(kj::ArrayPtr> pieces) override { if (writer == kj::none) { return KJ_EXCEPTION(FAILED, "Write after stream has been closed."); } auto amount = 0; for (auto& piece: pieces) { amount += piece.size(); } if (amount == 0) return kj::READY_NOW; return canceler.wrap(context.run([this, amount, pieces](Worker::Lock& lock) mutable { auto& writer = getInner(); // Sadly, we have to allocate and copy here. Our received set of buffers are only // guaranteed to live until the returned promise is resolved, but the application code // may hold onto the ArrayBuffer for longer. We need to make sure that the backing store // for the ArrayBuffer remains valid. auto source = KJ_ASSERT_NONNULL(jsg::BufferSource::tryAlloc(lock, amount)); auto ptr = source.asArrayPtr(); for (auto& piece: pieces) { KJ_DASSERT(ptr.size() > 0); KJ_DASSERT(piece.size() <= ptr.size()); if (piece.size() == 0) continue; ptr.first(piece.size()).copyFrom(piece); ptr = ptr.slice(piece.size()); } return context.awaitJs(lock, writer.write(lock, source.getHandle(lock))); })); } // TODO(perf): We can't properly implement tryPumpFrom(), which means that Cap'n Proto will // be unable to perform path shortening if the underlying stream turns out to be another capnp // stream. This isn't a huge deal, but might be nice to enable someday. It may require // significant refactoring of streams. kj::Promise whenWriteDisconnected() override { // TODO(soon): We might be able to support this by following the writer.closed promise, // which becomes resolved when the writer is used to close the stream, or rejects when // the stream has errored. However, currently, we don't have an easy way to do this. // // The Writer's getClosed() method returns a jsg::MemoizedIdentity>. // jsg::MemoizedIdentity lazily converts the jsg::Promise into a v8::Promise once it // passes through the type wrapper. It does not give us any way to consistently get // at the underlying jsg::Promise or the mapped v8::Promise. We would need to // capture a TypeHandler in here and convert each time to one or the other, then // attach our continuation. It's doable but a bit of a pain. // // For now, let's handle this the same as WritableStreamRpcAdapter and just return a // never done. return kj::NEVER_DONE; } kj::Promise end() override { if (writer == kj::none) { return KJ_EXCEPTION(FAILED, "End after stream has been closed."); } ended = true; return canceler.wrap(context.run([this](Worker::Lock& lock) mutable { return context.awaitJs(lock, getInner().close(lock)); })); } private: IoContext& context; kj::Maybe> writer; kj::Canceler canceler; kj::Own> doneFulfiller; kj::Own> weakRef = kj::refcounted>( kj::Badge(), *this); bool ended = false; WritableStreamDefaultWriter& getInner() { KJ_IF_SOME(inner, writer) { return *inner; } kj::throwFatalException(cancellationException()); } static kj::Exception cancellationException() { return JSG_KJ_EXCEPTION(DISCONNECTED, Error, "WritableStream received over RPC was disconnected because the remote execution context " "has endeded."); } }; } // namespace void WritableStream::serialize(jsg::Lock& js, jsg::Serializer& serializer) { // Serialize by effectively creating a `JsRpcStub` around this object and serializing that. // Except we don't actually want to do _exactly_ that, because we do not want to actually create // a `JsRpcStub` locally. So do the important parts of `JsRpcStub::constructor()` followed by // `JsRpcStub::serialize()`. auto& handler = JSG_REQUIRE_NONNULL(serializer.getExternalHandler(), DOMDataCloneError, "WritableStream can only be serialized for RPC."); auto externalHandler = dynamic_cast(&handler); JSG_REQUIRE(externalHandler != nullptr, DOMDataCloneError, "WritableStream can only be serialized for RPC."); IoContext& ioctx = IoContext::current(); // TODO(soon): Support JS-backed WritableStreams. Currently this only supports native streams // and IdentityTransformStream, since only they are backed by WritableStreamSink. KJ_IF_SOME(sink, getController().removeSink(js)) { // NOTE: We're counting on `removeSink()`, to check that the stream is not locked and other // common checks. It's important we don't modify the WritableStream before this call. auto encoding = sink->disownEncodingResponsibility(); auto wrapper = kj::heap(kj::mv(sink)); // Make sure this stream will be revoked if the IoContext ends. ioctx.addTask(wrapper->waitForCompletionOrRevoke().attach(ioctx.registerPendingEvent())); auto capnpStream = ioctx.getByteStreamFactory().kjToCapnp(kj::mv(wrapper)); externalHandler->write([capnpStream = kj::mv(capnpStream), encoding]( rpc::JsValue::External::Builder builder) mutable { auto ws = builder.initWritableStream(); ws.setByteStream(kj::mv(capnpStream)); ws.setEncoding(encoding); }); } else { // TODO(soon): Support disownEncodingResponsibility with JS-backed streams // NOTE: We're counting on `getWriter()` to check that the stream is not locked and other // common checks. It's important we don't modify the WritableStream before this call. auto wrapper = kj::heap(ioctx, getWriter(js)); // Make sure this stream will be revoked if the IoContext ends. ioctx.addTask(wrapper->waitForCompletionOrRevoke().attach(ioctx.registerPendingEvent())); auto capnpStream = ioctx.getByteStreamFactory().kjToCapnp(kj::mv(wrapper)); externalHandler->write( [capnpStream = kj::mv(capnpStream)](rpc::JsValue::External::Builder builder) mutable { auto ws = builder.initWritableStream(); ws.setByteStream(kj::mv(capnpStream)); ws.setEncoding(StreamEncoding::IDENTITY); }); } } jsg::Ref WritableStream::deserialize( jsg::Lock& js, rpc::SerializationTag tag, jsg::Deserializer& deserializer) { auto& handler = KJ_REQUIRE_NONNULL( deserializer.getExternalHandler(), "got WritableStream on non-RPC serialized object?"); auto externalHandler = dynamic_cast(&handler); KJ_REQUIRE(externalHandler != nullptr, "got WritableStream on non-RPC serialized object?"); auto reader = externalHandler->read(); KJ_REQUIRE(reader.isWritableStream(), "external table slot type doesn't match serialization tag"); auto ws = reader.getWritableStream(); auto encoding = ws.getEncoding(); KJ_REQUIRE( static_cast(encoding) < capnp::Schema::from().getEnumerants().size(), "unknown StreamEncoding received from peer"); IoContext& ioctx = IoContext::current(); auto stream = ioctx.getByteStreamFactory().capnpToKjExplicitEnd(ws.getByteStream()); auto sink = newSystemStream(kj::mv(stream), encoding, ioctx); return js.alloc( ioctx, kj::mv(sink), ioctx.getMetrics().tryCreateWritableByteStreamObserver()); } void WritableStreamDefaultWriter::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { KJ_IF_SOME(attached, state.tryGetActiveUnsafe()) { tracker.trackField("attached", attached.stream); } tracker.trackField("closedPromise", closedPromise); tracker.trackField("readyPromise", readyPromise); } void WritableStream::visitForMemoryInfo(jsg::MemoryTracker& tracker) const { tracker.trackField("controller", controller); } } // namespace workerd::api