#include "readable-source-adapter.h" #include "writable-sink.h" #include #include namespace workerd::api::streams { namespace { // Per the ReadableStream spec, when a read(buf) is performed on a BYOB reader, // if the stream is already closed, we still need to return the allocated buffer // back to the caller, but it must be in a zero-length view. This utility function // does that. It takes the original allocation and wraps it into a new ArrayBuffer // instance that is wrapped by a zero-length view of the same type as the original // TypedArray we were given. jsg::BufferSource transferToEmptyBuffer(jsg::Lock& js, jsg::BufferSource buffer) { KJ_DASSERT(!buffer.isDetached() && buffer.canDetach(js)); auto backing = buffer.detach(js); backing.limit(0); auto buf = jsg::BufferSource(js, kj::mv(backing)); KJ_DASSERT(buf.size() == 0); return kj::mv(buf); } } // namespace // The Active state maintains a queue of tasks, such as read or close operations. Each task // contains a promise-returning function object and a fulfiller. When the first task is // enqueued, the active state begins processing the queue asynchronously. Each function // is invoked in order, its promise awaited, and the result passed to the fulfiller. The // fulfiller notifies the code which enqueued the task that the task has completed. In // this way, read and close operations are safely executed in serial, even if one operation // is called before the previous completes. This mechanism satisfies KJ's restriction on // concurrent operations on streams. struct ReadableStreamSourceJsAdapter::Active { struct Task { kj::Function()> task; kj::Own> fulfiller; Task(kj::Function()> task, kj::Own> fulfiller) : task(kj::mv(task)), fulfiller(kj::mv(fulfiller)) {} KJ_DISALLOW_COPY_AND_MOVE(Task); }; using TaskQueue = workerd::util::Queue>; kj::Own source; kj::Canceler canceler; TaskQueue queue; bool canceled = false; bool running = false; bool closePending = false; kj::Maybe pendingCancel; Active(kj::Own source): source(kj::mv(source)) {} KJ_DISALLOW_COPY_AND_MOVE(Active); ~Active() noexcept(false) { // When the Active is dropped, we cancel any remaining pending reads and // abort the sink. cancel(KJ_EXCEPTION(DISCONNECTED, "Writable stream is canceled or closed.")); // Check invariants for safety. // 1. Our canceler should be empty because we canceled it. KJ_DASSERT(canceler.isEmpty()); // 2. The write queue should be empty. KJ_DASSERT(queue.empty()); } // Explicitly cancel all in-flight and pending tasks in the queue. // This is a non-op if cancel has already been called. void cancel(kj::Exception&& exception) { if (canceled) return; canceled = true; // 1. Cancel our in-flight "runLoop", if any. pendingCancel = exception.clone(); canceler.cancel(exception.clone()); // 2. Drop our queue of pending tasks. queue.drainTo( [&exception](kj::Own&& task) { task->fulfiller->reject(exception.clone()); }); // 3. Cancel and drop the source itself. We're done with it. if (exception.getType() != kj::Exception::Type::DISCONNECTED) { source->cancel(kj::mv(exception)); } auto dropped KJ_UNUSED = kj::mv(source); } kj::Promise enqueue(kj::Function()> task) { KJ_DASSERT(!canceled, "cannot enqueue tasks on a canceled queue"); auto paf = kj::newPromiseAndFulfiller(); queue.push(kj::heap(kj::mv(task), kj::mv(paf.fulfiller))); if (!running) { IoContext::current().addTask(canceler.wrap(run())); } return kj::mv(paf.promise); } kj::Promise run() { KJ_DEFER(running = false); running = true; while (!queue.empty() && !canceled) { auto task = KJ_ASSERT_NONNULL(queue.pop()); KJ_DEFER({ if (task->fulfiller->isWaiting()) { KJ_IF_SOME(pending, pendingCancel) { task->fulfiller->reject(kj::mv(pending)); } else { task->fulfiller->reject(KJ_EXCEPTION(DISCONNECTED, "Task was canceled.")); } } }); bool taskFailed = false; try { task->fulfiller->fulfill(co_await task->task()); } catch (...) { auto ex = kj::getCaughtExceptionAsKj(); task->fulfiller->reject(kj::mv(ex)); taskFailed = true; } // If the task failed, we exit the loop. We're going to abort the // entire remaining queue anyway so there's no point in continuing. if (taskFailed) co_return; } } }; ReadableStreamSourceJsAdapter::ReadableStreamSourceJsAdapter( jsg::Lock& js, IoContext& ioContext, kj::Own source) : state(State::create(ioContext.addObject(kj::heap(kj::mv(source))))), selfRef(kj::rc>( kj::Badge{}, *this)) {} ReadableStreamSourceJsAdapter::~ReadableStreamSourceJsAdapter() noexcept(false) { selfRef->invalidate(); } void ReadableStreamSourceJsAdapter::cancel(kj::Exception exception) { KJ_IF_SOME(open, state.tryGetActiveUnsafe()) { open.active->cancel(exception.clone()); } state.forceTransitionTo(kj::mv(exception)); } void ReadableStreamSourceJsAdapter::cancel(jsg::Lock& js, const jsg::JsValue& reason) { cancel(js.exceptionToKj(reason)); } void ReadableStreamSourceJsAdapter::shutdown(jsg::Lock& js) { KJ_IF_SOME(open, state.tryGetActiveUnsafe()) { open.active->cancel(KJ_EXCEPTION(DISCONNECTED, "Stream was shut down.")); state.transitionTo(); } // If we are are already closed or canceled, this is a no-op. } bool ReadableStreamSourceJsAdapter::isClosed() { return state.is(); } kj::Maybe ReadableStreamSourceJsAdapter::isCanceled() { return state.tryGetErrorUnsafe(); } jsg::Promise ReadableStreamSourceJsAdapter::read( jsg::Lock& js, ReadOptions options) { KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) { // Really should not have been called if errored but just in case, // return a rejected promise. return js.rejectedPromise(js.exceptionToJs(exception.clone())); } if (state.is()) { // We are already in a closed state. This is a no-op, just return // an empty buffer. return js.resolvedPromise(ReadResult{ .buffer = transferToEmptyBuffer(js, kj::mv(options.buffer)), .done = true, }); } auto& open = state.requireActiveUnsafe(); // Deference the IoOwn once to get the active state. Active& active = *open.active; // If close is pending, we cannot accept any more reads. // Treat them as if the stream is closed. if (active.closePending) { return js.resolvedPromise(ReadResult{ .buffer = transferToEmptyBuffer(js, kj::mv(options.buffer)), .done = true, }); } // Ok, we are in a readable state, there are no pending closes. // Let's enqueue our read request. auto& ioContext = IoContext::current(); auto buffer = kj::mv(options.buffer); auto elementSize = buffer.getElementSize(); // The buffer size should always be a multiple of the element size and should // always be at least as large as minBytes. This should be handled for us by // the jsg::BufferSource, but just to be safe, we will double-check with a // debug assert here. KJ_DASSERT(buffer.size() % elementSize == 0); auto minBytes = kj::min(options.minBytes.orDefault(elementSize), buffer.size()); // We want to be sure that minBytes is a multiple of the element size // of the buffer, otherwise we might never be able to satisfy the request // correcty. If the caller provided a minBytes, and it is not a multiple // of the element size, we will round it up to the next multiple. if (elementSize > 1) { minBytes = minBytes + (elementSize - (minBytes % elementSize)) % elementSize; } // Note: We do not enforce that the source must provide at least minBytes // if available here as that is part of the contract of the source itself. // We will simply pass minBytes along to the source and it is up to the // source to honor it. We do, however, enforce that the source must // never return more than the size of the buffer we provided. // We only pass a kj::ArrayPtr to the buffer into the read call, keeping // the actual buffer instance alive by attaching it to the JS promise // chain that follows the read in order to keep it alive. auto promise = active.enqueue(kj::coCapture( [&active, buffer = buffer.asArrayPtr(), minBytes]() mutable -> kj::Promise { // TODO(soon): The underlying kj streams API now supports passing the // kj::ArrayPtr directly to the read call, but ReadableStreamSource has // not yet been updated to do so. When it is, we can update this read to // pass `buffer` directly rather than passing the begin() and size(). co_return co_await active.source->read(buffer, minBytes); })); return ioContext .awaitIo(js, kj::mv(promise), [buffer = kj::mv(buffer), self = selfRef.addRef()](jsg::Lock& js, size_t bytesRead) mutable -> jsg::Promise { // If the bytesRead is 0, that indicates the stream is closed. We will // move the stream to a closed state and return the empty buffer. if (bytesRead == 0) { self->runIfAlive([](ReadableStreamSourceJsAdapter& self) { KJ_IF_SOME(open, self.state.tryGetActiveUnsafe()) { open.active->closePending = true; } }); return js.resolvedPromise(ReadResult{ .buffer = transferToEmptyBuffer(js, kj::mv(buffer)), .done = true, }); } KJ_DASSERT(bytesRead <= buffer.size()); // If bytesRead is not a multiple of the element size, that indicates // that the source either read less than minBytes (and ended), or is // simply unable to satisfy the element size requirement. We cannot // provide a partial element to the caller, so reject the read. if (bytesRead % buffer.getElementSize() != 0) { return js.rejectedPromise( js.typeError(kj::str("The underlying stream failed to provide a multiple of the " "target element size ", buffer.getElementSize()))); } auto backing = buffer.detach(js); backing.limit(bytesRead); return js.resolvedPromise(ReadResult{ .buffer = jsg::BufferSource(js, kj::mv(backing)), .done = false, }); }) .catch_(js, [self = selfRef.addRef()]( jsg::Lock& js, jsg::Value exception) -> ReadableStreamSourceJsAdapter::ReadResult { // If an error occurred while reading, we need to transition the adapter // to the canceled state, but only if the adapter is still alive. auto error = jsg::JsValue(exception.getHandle(js)); self->runIfAlive([&](ReadableStreamSourceJsAdapter& self) { self.cancel(js, error); }); js.throwException(kj::mv(exception)); }); } // Transitions the adapter into the closing state. Once the read queue // is empty, we will close the source and transition to the closed state. jsg::Promise ReadableStreamSourceJsAdapter::close(jsg::Lock& js) { KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) { // Really should not have been called if errored but just in case, // return a rejected promise. return js.rejectedPromise(js.exceptionToJs(exception.clone())); } if (state.is()) { // We are already in a closed state. This is a no-op. This really // should not have been called if closed but just in case, return // a resolved promise. return js.resolvedPromise(); } auto& open = state.requireActiveUnsafe(); auto& ioContext = IoContext::current(); auto& active = *open.active; if (active.closePending) { return js.rejectedPromise(js.typeError("Close already pending, cannot close again.")); } active.closePending = true; auto promise = active.enqueue([]() -> kj::Promise { co_return 0; }); return ioContext .awaitIo(js, kj::mv(promise), [self = selfRef.addRef()](jsg::Lock&, size_t) { self->runIfAlive( [](ReadableStreamSourceJsAdapter& self) { self.state.transitionTo(); }); }).catch_(js, [self = selfRef.addRef()](jsg::Lock& js, jsg::Value&& exception) { // Likewise, while nothing should be waiting on the ready promise, we // should still reject it just in case. auto error = jsg::JsValue(exception.getHandle(js)); self->runIfAlive([&](ReadableStreamSourceJsAdapter& self) { self.cancel(js, error); }); js.throwException(kj::mv(exception)); }); } jsg::Promise> ReadableStreamSourceJsAdapter::readAllText( jsg::Lock& js, uint64_t limit) { KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) { // Really should not have been called if errored but just in case, // return a rejected promise. return js.rejectedPromise>(js.exceptionToJs(exception.clone())); } if (state.is()) { // We are already in a closed state. This is a no-op. This really // should not have been called if closed but just in case, return // a resolved promise. return js.resolvedPromise(jsg::JsRef(js, js.str())); } auto& open = state.requireActiveUnsafe(); auto& ioContext = IoContext::current(); auto& active = *open.active; if (active.closePending) { return js.rejectedPromise>( js.typeError("Close already pending, cannot read.")); } active.closePending = true; struct Holder { kj::Maybe result; }; auto holder = kj::heap(); auto promise = active.enqueue([&active, &holder = *holder, limit]() -> kj::Promise { auto str = co_await active.source->readAllText(limit); size_t amount = str.size(); holder.result = kj::mv(str); co_return amount; }); return ioContext .awaitIo(js, kj::mv(promise), [self = selfRef.addRef(), holder = kj::mv(holder)](jsg::Lock& js, size_t amount) { self->runIfAlive( [&](ReadableStreamSourceJsAdapter& self) { self.state.transitionTo(); }); KJ_IF_SOME(result, holder->result) { KJ_DASSERT(result.size() == amount); return jsg::JsRef(js, js.str(result)); } else { return jsg::JsRef(js, js.str()); } }) .catch_(js, [self = selfRef.addRef()]( jsg::Lock& js, jsg::Value&& exception) -> jsg::JsRef { // Likewise, while nothing should be waiting on the ready promise, we // should still reject it just in case. auto error = jsg::JsValue(exception.getHandle(js)); self->runIfAlive([&](ReadableStreamSourceJsAdapter& self) { self.cancel(js, error); }); js.throwException(kj::mv(exception)); }); } jsg::Promise ReadableStreamSourceJsAdapter::readAllBytes( jsg::Lock& js, uint64_t limit) { KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) { // Really should not have been called if errored but just in case, // return a rejected promise. return js.rejectedPromise(js.exceptionToJs(exception.clone())); } if (state.is()) { // We are already in a closed state. This is a no-op. This really // should not have been called if closed but just in case, return // a resolved promise. auto backing = jsg::BackingStore::alloc(js, 0); return js.resolvedPromise(jsg::BufferSource(js, kj::mv(backing))); } auto& open = state.requireActiveUnsafe(); auto& ioContext = IoContext::current(); auto& active = *open.active; if (active.closePending) { return js.rejectedPromise( js.typeError("Close already pending, cannot read.")); } active.closePending = true; struct Holder { kj::Maybe> result; }; auto holder = kj::heap(); auto promise = active.enqueue([&active, &holder = *holder, limit]() -> kj::Promise { auto str = co_await active.source->readAllBytes(limit); size_t amount = str.size(); holder.result = kj::mv(str); co_return amount; }); return ioContext .awaitIo(js, kj::mv(promise), [self = selfRef.addRef(), holder = kj::mv(holder)](jsg::Lock& js, size_t amount) { self->runIfAlive( [&](ReadableStreamSourceJsAdapter& self) { self.state.transitionTo(); }); KJ_IF_SOME(result, holder->result) { KJ_DASSERT(result.size() == amount); // We have to copy the data into the backing store because of the // v8 sandboxing rules. auto backing = jsg::BackingStore::alloc(js, amount); backing.asArrayPtr().copyFrom(result); return jsg::BufferSource(js, kj::mv(backing)); } else { auto backing = jsg::BackingStore::alloc(js, 0); return jsg::BufferSource(js, kj::mv(backing)); } }) .catch_(js, [self = selfRef.addRef()](jsg::Lock& js, jsg::Value&& exception) -> jsg::BufferSource { // Likewise, while nothing should be waiting on the ready promise, we // should still reject it just in case. auto error = jsg::JsValue(exception.getHandle(js)); self->runIfAlive([&](ReadableStreamSourceJsAdapter& self) { self.cancel(js, error); }); js.throwException(kj::mv(exception)); }); } kj::Maybe ReadableStreamSourceJsAdapter::tryGetLength(StreamEncoding encoding) { KJ_IF_SOME(open, state.tryGetActiveUnsafe()) { return open.active->source->tryGetLength(encoding); } return kj::none; } kj::Maybe ReadableStreamSourceJsAdapter::tryTee( jsg::Lock& js, uint64_t limit) { KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) { js.throwException(js.exceptionToJs(exception.clone())); } if (state.is()) { // We are already closed, cannot tee. return kj::none; } auto& open = state.requireActiveUnsafe(); auto& active = *open.active; // If we are closing, or have pending tasks, we cannot tee. JSG_REQUIRE(!active.closePending && !active.running && active.queue.empty(), Error, "Cannot tee a stream that is closing or has pending reads."); auto tee = active.source->tee(limit); auto& ioContext = IoContext::current(); state.transitionTo(); return Tee{ .branch1 = kj::heap(js, ioContext, kj::mv(tee.branch1)), .branch2 = kj::heap(js, ioContext, kj::mv(tee.branch2)), }; } // =============================================================================================== struct ReadableSourceKjAdapter::Active { IoContext& ioContext; jsg::Ref stream; jsg::Ref reader; kj::Canceler canceler; struct Idle { static constexpr kj::StringPtr NAME KJ_UNUSED = "idle"_kj; }; struct Readable { static constexpr kj::StringPtr NAME KJ_UNUSED = "readable"_kj; // Previously read but unconsumed bytes. We keep these around for the next read call. kj::Array data; kj::ArrayPtr view; Readable(kj::Array&& data): data(kj::mv(data)), view(this->data) {} }; struct Reading { static constexpr kj::StringPtr NAME KJ_UNUSED = "reading"_kj; // The contract for ReadableStreamSource is that there can be only one read() in-flight // against the underlying stream at a time. }; struct Done { static constexpr kj::StringPtr NAME KJ_UNUSED = "done"_kj; // If a read returns fewer than the requested minBytes, that indicates the stream is done. We // make note of that here to prevent any further reads. We cannot transition to the closed // state in the promise chain of the read because the adapter will cancel the read promise // itself once Active is destroyed, and that would be a bad thing. }; struct Canceling { static constexpr kj::StringPtr NAME KJ_UNUSED = "canceling"_kj; kj::Exception exception; }; struct Canceled { static constexpr kj::StringPtr NAME KJ_UNUSED = "canceled"_kj; kj::Exception exception; }; // Inner state machine for tracking read operation state: // Idle -> Reading (start read) // Reading -> Idle (read complete, no leftover) // Reading -> Readable (read complete, has leftover) // Reading -> Done (read returned less than minBytes) // Any -> Canceling (error during read) // Any -> Canceled (explicit cancel) // Done, Canceling, and Canceled are terminal states. using InnerState = StateMachine, Idle, Readable, Reading, Done, Canceling, Canceled>; InnerState state; Active(jsg::Lock& js, IoContext& ioContext, jsg::Ref stream); KJ_DISALLOW_COPY_AND_MOVE(Active); ~Active() noexcept(false); void cancel(kj::Exception reason); }; // The ReadContext struct holds all the state needed to perform a read, // including the JS objects that need to be kept alive during the // read operation, the buffer we are reading into, and the total // number of bytes read so far. This must be kept alive until the // read is fully complete and returned back to the adapter when // the read is complete. // // Ownership of the ReadContext is passed into the isolate lock and // held by JS promise continuations, so it must not contain any // kj I/O objects or references without an IoOwn wrapper. struct ReadableSourceKjAdapter::ReadContext { jsg::Ref stream; jsg::Ref reader; kj::ArrayPtr buffer; // Only set to back the buffer if we need to keep it alive. kj::Maybe> backingBuffer; size_t totalRead = 0; size_t minBytes = 0; kj::Maybe maybeLeftOver; // We keep a weak reference to the adapter itself so we can track // whether it is still alive while we are in a JS promise chain. // If the adapter is gone, or transitions to a closed or canceled // state we will abandon the read. If the ref is not set, then we // are in a pump operation and do not need to check for liveness. kj::Maybe>> adapterRef; void reset() { // Resetting is only allowed if we have the backing buffer. buffer = KJ_ASSERT_NONNULL(backingBuffer); totalRead = 0; minBytes = 0; maybeLeftOver = kj::none; } }; namespace { constexpr size_t kMinRemainingForAdditionalRead = 512; jsg::Ref initReader(jsg::Lock& js, jsg::Ref& stream) { JSG_REQUIRE(!stream->isLocked(), TypeError, "ReadableStream is locked."); JSG_REQUIRE(!stream->isDisturbed(), TypeError, "ReadableStream is disturbed."); auto reader = stream->getReader(js, kj::none); return kj::mv(KJ_ASSERT_NONNULL(reader.tryGet>())); } using JsByteSource = kj::OneOf, jsg::JsRef, jsg::JsRef>; kj::Maybe tryExtractJsByteSource(jsg::Lock& js, const jsg::JsValue& jsval) { KJ_IF_SOME(abView, jsval.tryCast()) { return kj::Maybe(jsg::JsRef(js, abView)); } else KJ_IF_SOME(ab, jsval.tryCast()) { return kj::Maybe(jsg::JsRef(js, ab)); } else KJ_IF_SOME(str, jsval.tryCast()) { return kj::Maybe(jsg::JsRef(js, str)); } return kj::none; } // Copies as much data from source into the context as possible, returning // the number of bytes copied. kj::Maybe> copyFromSource( jsg::Lock& js, ReadableSourceKjAdapter::ReadContext& context, const JsByteSource& source) { KJ_SWITCH_ONEOF(source) { KJ_CASE_ONEOF(str, jsg::JsRef) { auto view = str.getHandle(js); size_t len = view.length(js); size_t toCopy = kj::min(len, context.buffer.size()); if (toCopy == 0) { return kj::none; } if (toCopy < len) { // We are going to have left-over data. Unfortunately in this case // we have to copy the data twice... once into a kj::String and // again into our buffer. This is because the V8 string UTF-8 // write API does not support partial writes with an offset. auto data = view.toUSVString(js); context.buffer.first(toCopy).copyFrom(data.asBytes().first(toCopy)); context.totalRead += toCopy; context.buffer = context.buffer.slice(toCopy); KJ_DASSERT(context.buffer.size() == 0); return kj::Maybe(data.asBytes().slice(toCopy).attach(kj::mv(data))); } // We can copy everything in one go. Yay! This is great because we // can avoid a double copy here. auto ret KJ_UNUSED = view.writeInto(js, context.buffer.asChars().first(toCopy), jsg::JsString::WriteFlags::REPLACE_INVALID_UTF8); KJ_DASSERT(ret.written == toCopy); context.totalRead += toCopy; context.buffer = context.buffer.slice(toCopy); return kj::none; } KJ_CASE_ONEOF(ab, jsg::JsRef) { auto src = ab.getHandle(js).asArrayPtr(); size_t toCopy = kj::min(src.size(), context.buffer.size()); if (toCopy == 0) { return kj::none; } context.buffer.first(toCopy).copyFrom(src.first(toCopy)); context.totalRead += toCopy; context.buffer = context.buffer.slice(toCopy); if (toCopy < src.size()) { KJ_DASSERT(context.buffer.size() == 0); // TODO(mpk): For now, we have to copy the left-over data into a new array. // Why? I'm happy you asked! Because the src is backed by a // v8::BackingStore protected by the v8 sandboxing rules and we // don't yet have the memory protection key logic in place to safely // share that memory outside of the v8 heap. For now, copy. Later // we can revisit this to hopefully avoid the additinal copy. return kj::Maybe(kj::heapArray(src.slice(toCopy))); } return kj::none; } KJ_CASE_ONEOF(view, jsg::JsRef) { auto src = view.getHandle(js).asArrayPtr(); size_t toCopy = kj::min(src.size(), context.buffer.size()); if (toCopy == 0) { // Copy nothing. Return 0. return kj::none; } context.buffer.first(toCopy).copyFrom(src.first(toCopy)); context.totalRead += toCopy; context.buffer = context.buffer.slice(toCopy); if (toCopy < src.size()) { KJ_DASSERT(context.buffer.size() == 0); return kj::Maybe(kj::heapArray(src.slice(toCopy))); } return kj::none; } } KJ_UNREACHABLE; } } // namespace ReadableSourceKjAdapter::Active::Active( jsg::Lock& js, IoContext& ioContext, jsg::Ref stream) : ioContext(ioContext), stream(kj::mv(stream)), reader(initReader(js, this->stream)), state(InnerState::create()) {} ReadableSourceKjAdapter::Active::~Active() noexcept(false) { cancel(KJ_EXCEPTION(DISCONNECTED, "ReadableSourceKjAdapter is canceled.")); } void ReadableSourceKjAdapter::Active::cancel(kj::Exception reason) { if (state.is()) { return; } bool wasDone = state.is(); state.forceTransitionTo(reason.clone()); canceler.cancel(reason.clone()); if (!wasDone) { // If the previous read indicated that it was the last read, then // the reader will have already been dropped. We do not need to // cancel it here. ioContext.addTask(ioContext.run([readable = kj::mv(stream), reader = kj::mv(reader), exception = kj::mv(reason)](jsg::Lock& js) mutable { auto& ioContext = IoContext::current(); auto error = js.exceptionToJsValue(kj::mv(exception)); auto promise = reader->cancel(js, error.getHandle(js)); return ioContext.awaitJs(js, kj::mv(promise)); })); } } ReadableSourceKjAdapter::ReadableSourceKjAdapter( jsg::Lock& js, IoContext& ioContext, jsg::Ref stream, Options options) : state(KjState::create(kj::heap(js, ioContext, kj::mv(stream)))), options(options), selfRef( kj::rc>(kj::Badge{}, *this)) {} ReadableSourceKjAdapter::~ReadableSourceKjAdapter() noexcept(false) { selfRef->invalidate(); } jsg::Promise> ReadableSourceKjAdapter::readInternal( jsg::Lock& js, kj::Own context, MinReadPolicy minReadPolicy) { auto& ioContext = IoContext::current(); // Pay close attention to the lambda captures here. There are no raw references // captured! The adapter itself may be destroyed or closed while we are in the // promise chain below, so we have to be careful to only hold weak references // and pass ownership of the context along the promise chain. // // The other important thing here is to remember that everything in this function // is running within the isolate lock. The idea is to keep the entire read of the // underlying stream entirely within the lock so that we don't have to bounce // in and out of the isolate lock multiple times. We only return to the kj world // once the entire read is complete. // // Note the uses of addFunctor below. This is important because it ensures // that the promise continuations are run within the correct IoContext. return context->reader->read(js).then(js, ioContext.addFunctor([context = kj::mv(context), minReadPolicy](jsg::Lock& js, ReadResult result) mutable -> jsg::Promise> { if (result.done || result.value == kj::none) { // Stream is ended. return js.resolvedPromise(kj::mv(context)); } auto& value = KJ_ASSERT_NONNULL(result.value); // Ok, we have some data. Let's make sure it is bytes. // We accept either an ArrayBuffer, ArrayBufferView, or string. auto jsval = jsg::JsValue(value.getHandle(js)); KJ_IF_SOME(result, tryExtractJsByteSource(js, jsval)) { // Process the resulting data. KJ_IF_SOME(leftOver, copyFromSource(js, *context, result)) { KJ_ASSERT(context->buffer.size() == 0); if (leftOver.size() > 0) { context->maybeLeftOver = Active::Readable(kj::mv(leftOver)); } else { context->maybeLeftOver = kj::none; } return js.resolvedPromise(kj::mv(context)); } // At this point, we should have no left over data. KJ_DASSERT(context->maybeLeftOver == kj::none); // If the buffer is exactly full (the chunk filled it perfectly), we're done. if (context->buffer.size() == 0) { return js.resolvedPromise(kj::mv(context)); } // We might continue reading only if the adapter is still alive and // in an active state... bool continueReading = true; KJ_IF_SOME(adapterRef, context->adapterRef) { continueReading = adapterRef->isValid(); adapterRef->runIfAlive( [&](ReadableSourceKjAdapter& adapter) { continueReading = adapter.state.isActive(); }); } // If we have satisfied the minimum read requirement and either // (a) the minReadPolicy is IMMEDIATE or (b) there are fewer // than 512 bytes left in the buffer, we will just return what we // have. The idea here is that while we could just return what we have // and let the caller call read again, that would be inefficient if // the caller has a large buffer and is trying to read a lot of data. // Instead of returning early with a minimally filled buffer, let's // try to fill it up a bit more before returning. The 512 byte limit // is somewhat arbitrary. The risk, of course, is that the next read // will return too much data to fit into the buffer, which will then // have to be stashed away as left over data. There's also a risk that // the stream is slow and we end up with more latency waiting for // the next chunk of data to arrive. In practice, this seems unlikely // to be a problem. The IMMEDIATE policy is useful in the latter case, // when the caller wants to get whatever data is available as soon // as possible, even if it is just a small amount. The downside of the // IMMEDIATE policy is that it can lead to a lot of small reads that // are expensive because they have to grab the isolate lock each time. bool minReadSatisfied = context->totalRead >= context->minBytes && (minReadPolicy == MinReadPolicy::IMMEDIATE || context->buffer.size() < kMinRemainingForAdditionalRead); if (!continueReading || minReadSatisfied) { return js.resolvedPromise(kj::mv(context)); } // We still have not satisfied the minimum read requirement or we are // trying to fill up a larger buffer. We will need to read more. Let's // call readInternal again to get the next chunk of data. Keep in mind // that this is not a true recursive call because readInternal returns // a jsg::Promise. We're just chaining the promises together here. return readInternal(js, kj::mv(context), minReadPolicy); } // Oooo, invalid type. We cannot handle this and must treat this as a fatal error. // We will cancel the stream and return an error. auto error = js.typeError("ReadableStream provided a non-bytes value. Only ArrayBuffer, " "ArrayBufferView, or string are supported."); context->reader->cancel(js, error); return js.rejectedPromise>(error); }), ioContext.addFunctor([](jsg::Lock& js, jsg::Value exception) { // In this case, the reader should already be in an errored state // since it it the read that failed. Just propagate the error. return js.rejectedPromise>(kj::mv(exception)); })); } // We separate out the actual read implementation so that it can be used by // both read and the pumpToImpl implementation. kj::Promise ReadableSourceKjAdapter::readImpl( Active& active, kj::ArrayPtr dest, size_t minBytes) { KJ_IF_SOME(readable, active.state.tryGetUnsafe()) { // We have some data left over from a previous read. Use that first. // If we have enough left over to fully satisfy this read, // Use it, then update our left over view. if (readable.view.size() >= dest.size()) { dest.copyFrom(readable.view.first(dest.size())); readable.view = readable.view.slice(dest.size()); if (readable.view.size() == 0) { // We used up all our left over data. We can transition to the idle state. active.state.transitionTo(); } // Otherwise we still have some left over data. That // is ok, we will keep it around for the next read. // We intentionally do not transition to the idle state // here because we want to keep the left over data for // the next read. return dest.size(); } // Otherwise, consume what we do have left over. auto size = readable.view.size(); dest.first(size).copyFrom(readable.view); dest = dest.slice(size); active.state.transitionTo(); // Did we at least satisfy the minimum bytes? if (size >= minBytes) { // Awesome, we are technically done with this read. // While we might actually have more room in our buffer, and the // minReadyPolicy might be OPPORTUNISTIC, we will not try to // read more from the stream right now so that we can avoid having // to grab the isolate lock for this read. Instead, let's return // what we have and let the caller call read again if/when they want. // This risks leaving a fair amount of unused space in the buffer // and requiring more read calls but it avoids the overhead of // an additional isolate lock grab when we know we can at least // provide some data right now. return size; } } // If we got here, we still have not satisfied the minimum bytes, // so we will continue on to read more from the stream. But, we // also should not have any more data left over. Let's verify. KJ_ASSERT(active.state.is()); active.state.transitionTo(); // Our read context holds all the state needed to perform the read. // Ownership of the context is passed into the read operation and // returned back to us when the read is complete. auto context = kj::heap({ .stream = active.stream.addRef(), .reader = active.reader.addRef(), .buffer = dest, .totalRead = 0, .minBytes = minBytes, .adapterRef = selfRef.addRef(), }); return active.canceler .wrap( // Warning: Do *not* capture "active" in this lambda! It may be destroyed // while we are in the promise chain. Instead, we capture a weak // reference to the adapter itself and check that we are still alive // and active before trying to update any state. active.ioContext.run([context = kj::mv(context), self = selfRef.addRef(), minReadPolicy = options.minReadPolicy]( jsg::Lock& js) mutable -> kj::Promise { auto& ioContext = IoContext::current(); // Perform the actual read. return ioContext.awaitJs(js, readInternal(js, kj::mv(context), minReadPolicy)) .then([self = kj::mv(self)](kj::Own context) mutable -> kj::Promise { // By the time we get here, it is possible that the adapter has been // destroyed. If that's the case, it's okay, that's what our weak ref // is here for. We will only try to update our state if we are still // alive and active. self->runIfAlive([&](ReadableSourceKjAdapter& self) { // Ok, we're still alive! Yay! But, let's check to make sure we didn't // change state while we were reading. KJ_IF_SOME(open, self.state.tryGetActiveUnsafe()) { auto& active = *open.active; // Ok, we're still active. Let's see if we have any left over data // that we need to stash away for the next read. KJ_IF_SOME(leftOver, context->maybeLeftOver) { // We have some left over data. Stash it away for the next read. active.state.transitionTo(kj::mv(leftOver)); // In this branch, we must have filled the entire destination // buffer and satisfied the minimum read requirement or else // we wouldn't have any left over data. Let's just assert that // invariant just in case. KJ_DASSERT(context->totalRead >= context->minBytes); } else if (context->totalRead < context->minBytes) { // We returned fewer than the minimum bytes requested. This is our // signal that we're done. active.state.transitionTo(); // We cannot change the state to Closed here because we are still // inside the kj::Promise chain wrapped by the canceler. If we // change the state to Closed, the Active would be destroyed, causing // this promise chain to be canceled. auto droppedReader KJ_UNUSED = kj::mv(active.reader); auto droppedStream KJ_UNUSED = kj::mv(active.stream); // In this branch, we should not have any left over data. // Let's assert that invariant just in case. KJ_DASSERT(context->maybeLeftOver == kj::none); } else { // Our read is complete. Return to the idle state and we're done. active.state.transitionTo(); // In this branch, we must have satisfied the minimum read // requirement. Let's just assert that invariant just in case. KJ_DASSERT(context->totalRead >= context->minBytes); // We should not have any left over data. KJ_DASSERT(context->maybeLeftOver == kj::none); } } else { // We were closed or canceled while we were reading. Doh! // That's ok, there's nothing more we can or need to do // here. Just fall-through to the return below. } }); return context->totalRead; }); })).catch_([self = selfRef.addRef()](kj::Exception exception) -> kj::Promise { self->runIfAlive([&](ReadableSourceKjAdapter& self) { KJ_IF_SOME(open, self.state.tryGetActiveUnsafe()) { open.active->state.forceTransitionTo(Active::Canceling{ .exception = exception.clone(), }); } }); return kj::mv(exception); }); } kj::Promise ReadableSourceKjAdapter::read(kj::ArrayPtr buffer, size_t minBytes) { if (buffer.size() == 0) { // Nothing to read. This is a no-op. return static_cast(0); } // Clamp the minBytes to [1, buffer.size()]. minBytes = kj::min(buffer.size(), kj::max(minBytes, 1UL)); KJ_DASSERT(minBytes >= 1 && minBytes <= buffer.size(), "minBytes must be less than or equal to the buffer size."); KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) { return exception.clone(); } if (state.is()) { return static_cast(0); } auto& open = state.requireActiveUnsafe(); auto& active = *open.active; KJ_SWITCH_ONEOF(active.state) { KJ_CASE_ONEOF(_, Active::Reading) { KJ_FAIL_REQUIRE("Cannot have multiple concurrent reads."); } KJ_CASE_ONEOF(_, Active::Done) { // The previous read indicated that it was the last read by returning // less than the minimum bytes requested. We have to treat this as // the stream being closed. state.transitionTo(); return static_cast(0); } KJ_CASE_ONEOF(canceling, Active::Canceling) { // The stream is being canceled. Propagate the exception and complete // the state transition. return KJ_ASSERT_NONNULL(checkCancelingOrCanceled(active)); } KJ_CASE_ONEOF(canceled, Active::Canceled) { // The stream was canceled. Propagate the exception and complete // the state transition. return KJ_ASSERT_NONNULL(checkCancelingOrCanceled(active)); } KJ_CASE_ONEOF(r, Active::Readable) { // There is some data left over from a previous read. return readImpl(active, buffer, minBytes); } KJ_CASE_ONEOF(_, Active::Idle) { // There are no pending reads and no left over data. return readImpl(active, buffer, minBytes); } } KJ_UNREACHABLE; } kj::Maybe ReadableSourceKjAdapter::tryGetLength(StreamEncoding encoding) { KJ_IF_SOME(open, state.tryGetActiveUnsafe()) { auto& active = *open.active; if (active.state.is() || active.state.is()) { // If the previous read indicated that it was the last, then // let's just transition to the closed state now and return kj::none. state.transitionTo(); return kj::none; } if (checkCancelingOrCanceled(active) != kj::none) { return kj::none; } return active.stream->tryGetLength(encoding).map( [](uint64_t len) { return static_cast(len); }); } // The stream is either closed or errored. return kj::none; } void ReadableSourceKjAdapter::cancel(kj::Exception reason) { KJ_IF_SOME(open, state.tryGetActiveUnsafe()) { open.active->cancel(reason.clone()); } state.forceTransitionTo(kj::mv(reason)); } kj::Maybe ReadableSourceKjAdapter::checkCancelingOrCanceled(Active& active) { KJ_IF_SOME(canceling, active.state.tryGetUnsafe()) { auto exception = kj::mv(canceling.exception); state.forceTransitionTo(exception.clone()); return kj::mv(exception); } KJ_IF_SOME(canceled, active.state.tryGetUnsafe()) { auto exception = kj::mv(canceled.exception); state.forceTransitionTo(exception.clone()); return kj::mv(exception); } return kj::none; } void ReadableSourceKjAdapter::throwIfCancelingOrCanceled(Active& active) { KJ_IF_SOME(exception, checkCancelingOrCanceled(active)) { kj::throwFatalException(kj::mv(exception)); } } kj::Promise ReadableSourceKjAdapter::pumpToImpl( kj::Own active, WritableSink& output, EndAfterPump end) { // This implementation uses DrainingReader to efficiently pull all synchronously // available data from the underlying JS stream in each iteration. This minimizes // the number of isolate lock acquisitions by getting all available data at once // rather than reading into fixed-size buffers. KJ_DASSERT(active->state.is() || active->state.is(), "pumpToImpl called when stream is not in an active state."); bool writeFailed = false; // First, if the active state is in the Readable state, we need to drain the // left over data before starting the main read loop. // This is unlikely to occur in the typical case, but we need to handle it // nonetheless. KJ_IF_SOME(readable, active->state.tryGetUnsafe()) { co_await output.write(readable.view); active->state.transitionTo(); } // We hold the DrainingReader during the pump. The pointer remains valid because // the reader is created and owned during the pump loop lifetime. kj::Maybe> maybeReader; // Initialize the pump by releasing the default reader and creating a DrainingReader. // This requires the isolate lock. co_await active->ioContext.run( [&active, &maybeReader](jsg::Lock& js) mutable -> kj::Promise { // Release the existing reader's lock so we can create a DrainingReader. active->reader->releaseLock(js); // Create the DrainingReader for the stream. maybeReader = KJ_ASSERT_NONNULL(DrainingReader::create(js, *active->stream), "Failed to create DrainingReader - stream should not be locked"); return kj::READY_NOW; }); auto& reader = KJ_ASSERT_NONNULL(maybeReader); kj::Maybe pendingException; try { while (true) { // Perform a draining read to get all synchronously available data. // Pass raw pointer to reader into the lambda - safe because we own it // and keep it alive for the duration of the pump. // The draining reader grabs all data currently available in the stream's // queue, then tries to read more data up to a limit as long as the data // can be provided synchronously. The idea is to drain off as much data // from the stream as possible each time we are holding the isolate lock // to minimize the number of times we need to re-enter the lock. DrainingReader* readerPtr = reader.get(); DrainingReadResult result = co_await active->ioContext.run([readerPtr](jsg::Lock& js) mutable { auto& ioContext = IoContext::current(); // Use a 256KB limit to allow periodic yielding to the event loop, // preventing a fast producer from monopolizing the thread. This limit // only affects subsequent pump iterations after the initial buffer drain. constexpr size_t kMaxReadPerCycle = 256 * 1024; return ioContext.awaitJs(js, readerPtr->read(js, kMaxReadPerCycle)); }); // Write all the chunks we received using vectored write for efficiency. if (result.chunks.size() > 0) { KJ_ON_SCOPE_FAILURE(writeFailed = true); // Convert Array> to ArrayPtr> for vectored write. auto pieces = KJ_MAP(chunk, result.chunks) -> kj::ArrayPtr { return chunk.asPtr(); }; co_await output.write(pieces); } // If the stream is done, end the output if needed and exit. if (result.done) { KJ_ON_SCOPE_FAILURE(writeFailed = true); if (end) { co_await output.end(); } co_return; } } } catch (...) { auto exception = kj::getCaughtExceptionAsKj(); if (!writeFailed) { // If we got an error and it wasn't the write that failed, abort the output. output.abort(exception.clone()); } // Store the exception to handle after the catch block. pendingException = kj::mv(exception); } // If there was an error, cancel the reader and propagate the exception. KJ_IF_SOME(exception, pendingException) { DrainingReader* readerPtr = reader.get(); co_await active->ioContext.run([readerPtr, ex = exception.clone()](jsg::Lock& js) mutable { auto& ioContext = IoContext::current(); auto error = js.exceptionToJsValue(kj::mv(ex)); return ioContext.awaitJs(js, readerPtr->cancel(js, error.getHandle(js))); }); kj::throwFatalException(kj::mv(exception)); } } kj::Promise> ReadableSourceKjAdapter::pumpTo( WritableSink& output, EndAfterPump end) { // The pumpTo operation continually reads from the stream and writes // to the output until the stream is closed or an error occurs. Once // the pump starts, the adapter transitions to the closed state and // ownership of the underlying stream is transferred to the pump // operation. KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) { return kj::Promise>(DeferredProxy{exception.clone()}); } if (state.is()) { // Already closed, nothing to do. return newNoopDeferredProxy(); } auto& open = state.requireActiveUnsafe(); auto& active = *open.active; // Per the contract for ReadableStreamSource::pumpTo, the pump operation // will take over ownership of the underlying stream until it is complete, // leaving the adapter itself in a closed state once the pump starts. // Dropping the returned promise will cancel the pump operation. // We do, however, need to first make sure that our active state is // not already pending a read or terminal state change. KJ_REQUIRE(!active.state.is(), "Cannot have multiple concurrent reads."); if (active.state.is()) { // The previous read indicated that it was the last read by returning // less than the minimum bytes requested, or the stream was fully // canceled. We have to treat this as the stream being closed. state.transitionTo(); return newNoopDeferredProxy(); } KJ_IF_SOME(exception, checkCancelingOrCanceled(active)) { return kj::Promise>(kj::mv(exception)); } // The active state should be Readable of Idle here. Let's verify. KJ_DASSERT(active.state.is() || active.state.is()); // The Active state will be transferred into the pumpImpl operation. auto activeState = kj::mv(open.active); state.transitionTo(); // transition to closed immediately // Because pumpToImpl is wrapping a JavaScript stream, it is not eligible // for deferred proxying. We will return a noopDeferredProxy that wraps the // promise from pumpToImpl(); return addNoopDeferredProxy(pumpToImpl(kj::mv(activeState), output, end)); } ReadableSource::Tee ReadableSourceKjAdapter::tee(size_t) { KJ_UNIMPLEMENTED("Teeing a ReadableSourceKjAdapter is not supported."); // Explanation: Teeing a ReadableStream must be done under the isolate lock, // as does creating a new ReadableSourceKjAdapter. However, when tee() // is called we are not guaranteed to be under the isolate lock, nor can // we acquire the lock here because this is a synchronous operation and // acquiring the isolate lock requires waiting for a promise to resolve. // // Teeing here is unlikely to be necessary. If you do need a tee, it's // necessary to tee the underlying ReadableStream directly and create // two separate ReadableSourceKjAdapters, one for each branch of // that tee while the lock is held. } kj::Promise> ReadableSourceKjAdapter::readAllBytes(size_t limit) { co_return co_await readAllImpl(limit); } kj::Promise ReadableSourceKjAdapter::readAllText(size_t limit) { auto array = co_await readAllImpl(limit); co_return kj::String(kj::mv(array)); } template kj::Promise> ReadableSourceKjAdapter::readAllImpl(size_t limit) { KJ_IF_SOME(exception, state.tryGetErrorUnsafe()) { kj::throwFatalException(exception.clone()); } if (state.is()) { co_return kj::Array(); } auto& open = state.requireActiveUnsafe(); auto& active = *open.active; KJ_REQUIRE(!active.state.is(), "Cannot have multiple concurrent reads."); if (active.state.is()) { // The previous read indicated that it was the last read by returning // less than the minimum bytes requested. We have to treat this as // the stream being closed. state.transitionTo(); co_return kj::Array(); } throwIfCancelingOrCanceled(active); // Our readAll operation will accumulate data into a buffer up to the // specified limit. If the limit is exceeded, the returned promise will // be rejected. Once the readAll operation starts, the adapter is moved // into a closed state and ownership of the underlying stream is transferred // to the readAll promise. auto activeState = kj::mv(open.active); state.transitionTo(); // transition to closed immediately KJ_DASSERT(activeState->state.is() || activeState->state.is()); // We do not use the canceler here. The adapter is closed and can be safely dropped. // This promise, however, will keep the stream alive until the read is completed. // If the returned promise is dropped, the readAll operation will be canceled. CancelationToken cancelationToken; co_return co_await IoContext::current().run( [limit, active = kj::mv(activeState), cancelationToken = cancelationToken.getWeakRef()]( jsg::Lock& js) mutable -> kj::Promise> { kj::Vector accumulated; // If we know the length of the stream ahead of time, and it is within the limit, // we can reserve that much space in the accumulator to avoid multiple allocations. KJ_IF_SOME(length, active->stream->tryGetLength(StreamEncoding::IDENTITY)) { if (length <= limit) { accumulated.reserve(length); // Pre-allocate } } auto& ioContext = IoContext::current(); return ioContext.awaitJs(js, readAllReadImpl(js, ioContext.addObject(kj::mv(active)), kj::mv(accumulated), limit, kj::mv(cancelationToken))); }); } template jsg::Promise> ReadableSourceKjAdapter::readAllReadImpl(jsg::Lock& js, IoOwn active, kj::Vector accumulated, size_t limit, kj::Rc> cancelationToken) { // Check for cancelation. The cancelation token is a weak ref. If the promise // that represents the readAll operation is dropped, the token will be invalidated. // Since there is no way to directly cancel a JavaScript promise, this is the best // we can do to interrupt the loop. if (!cancelationToken->isValid()) { return js.rejectedPromise>(js.error("readAll operation was canceled.")); } // First, drain any leftover data if the active state is in Readable mode. KJ_IF_SOME(readable, active->state.tryGetUnsafe()) { auto leftover = readable.view.asBytes(); if (leftover.size() > limit) { auto error = js.rangeError("Memory limit would be exceeded before EOF."); return active->reader->cancel(js, error).then( js, [ex = jsg::JsRef(js, error)](jsg::Lock& js) { return js.rejectedPromise>(ex.getHandle(js)); }); } if constexpr (kj::isSameType()) { accumulated.addAll(leftover.asChars()); } else { accumulated.addAll(leftover); } active->state.transitionTo(); } return active->reader->read(js).then(js, [active = kj::mv(active), accumulated = kj::mv(accumulated), limit, cancelationToken = kj::mv(cancelationToken)]( jsg::Lock& js, ReadResult result) mutable -> jsg::Promise> { // Check for cancelation. if (!cancelationToken->isValid()) { return js.rejectedPromise>(js.error("readAll operation was canceled.")); } if (result.done || result.value == kj::none) { // Stream ended. Return accumulated data. // If we're reading text, add NUL terminator. if constexpr (kj::isSameType()) { accumulated.add('\0'); } return js.resolvedPromise(accumulated.releaseAsArray()); } auto& value = KJ_ASSERT_NONNULL(result.value); auto jsval = jsg::JsValue(value.getHandle(js)); kj::ArrayPtr bytes; kj::Maybe maybeOwnedString; KJ_IF_SOME(str, jsval.tryCast()) { auto data = str.toUSVString(js); bytes = data.asBytes(); maybeOwnedString = kj::mv(data); } else KJ_IF_SOME(ab, jsval.tryCast()) { bytes = ab.asArrayPtr(); } else KJ_IF_SOME(view, jsval.tryCast()) { bytes = view.asArrayPtr(); } else { auto error = js.typeError("ReadableStream provided a non-bytes value. Only ArrayBuffer, " "ArrayBufferView, or string are supported."); return active->reader->cancel(js, error).then( js, [err = jsg::JsRef(js, error)](jsg::Lock& js) { return js.rejectedPromise>(err.getHandle(js)); }); } if (accumulated.size() + bytes.size() > limit) { auto error = js.rangeError("Memory limit would be exceeded before EOF."); return active->reader->cancel(js, error).then( js, [err = jsg::JsRef(js, error)](jsg::Lock& js) { return js.rejectedPromise>(err.getHandle(js)); }); } // Accumulate the bytes. if constexpr (kj::isSameType()) { accumulated.addAll(bytes.asChars()); } else { accumulated.addAll(bytes); } // Continue reading. return readAllReadImpl( js, kj::mv(active), kj::mv(accumulated), limit, kj::mv(cancelationToken)); }); } } // namespace workerd::api::streams