File
Blob: src/workerd/api/streams/readable-source.h
| 1 | #pragma once |
| 2 | |
| 3 | #include <workerd/io/worker-interface.capnp.h> |
| 4 | #include <workerd/util/strong-bool.h> |
| 5 | |
| 6 | #include <kj/debug.h> |
| 7 | |
| 8 | namespace kj { |
| 9 | class AsyncInputStream; |
| 10 | } |
| 11 | |
| 12 | namespace workerd { |
| 13 | |
| 14 | class IoContext; |
| 15 | namespace api { |
| 16 | template <typename T> |
| 17 | struct DeferredProxy; |
| 18 | class ReadableStreamSource; |
| 19 | } // namespace api |
| 20 | |
| 21 | namespace jsg { |
| 22 | class Lock; |
| 23 | } // namespace jsg |
| 24 | |
| 25 | namespace api::streams { |
| 26 | |
| 27 | class WritableSink; |
| 28 | |
| 29 | WD_STRONG_BOOL(EndAfterPump); |
| 30 | |
| 31 | // A ReadableSource is primarily intended to serve as a bridge between kj::AsyncInputStream |
| 32 | // and the ReadableStream API. However, it can also be used directly by KJ-space code that needs |
| 33 | // deferred proxying. While ReadableSource should probably have been a more JS-friendly |
| 34 | // API, it's a bit too late to change that now. Use the ReadableSourceJsAdapter in the |
| 35 | // readable-source-adapter.h file to wrap a ReadableSource for use from JavaScript. |
| 36 | // |
| 37 | // A ReadableSource must be treated like a KJ I/O object. Instances that are held |
| 38 | // by any JS-heap objects must be held by an IoOwn. |
| 39 | // |
| 40 | // If the ReadableSource is canceled or dropped, all pending read() reads will be |
| 41 | // canceled. |
| 42 | // |
| 43 | // Only one read() may be pending at a time. Attempting to initiate a second read() |
| 44 | // while one is already pending will result in a rejected promise. |
| 45 | // |
| 46 | // Calling pumpTo initiates a sequence of read() calls until the stream is fully consumed. |
| 47 | // Ownership of the underlying AsyncInputStream is transferred to the pumpTo operation and |
| 48 | // the ReadableSource is put into a closed state. After calling pumpTo, no further |
| 49 | // read() calls may be made directly on the ReadableSource. Dropping the returned |
| 50 | // promise before it resolves will cancel the pump operation. |
| 51 | // |
| 52 | // It is **NOT** intended that you should implement this interface for general use. |
| 53 | // It is only intended to be implemented by specific classes within workerd for the |
| 54 | // purpose of bridging between kj/js streams. Streams that operate at the kj level |
| 55 | // should implement the kj::Async*Stream interfaces, and streams that operate at |
| 56 | // the JS level should implement to the UnderlyingSource interface. This is a |
| 57 | // departure from what we've done previously but as part of the effort to simplify |
| 58 | // the streams code, the goal is to reduce the number of different stream interfaces |
| 59 | // that we implement to. |
| 60 | class ReadableSource { |
| 61 | public: |
| 62 | // Read into the given buffer, returning a promise that resolves to the number of bytes read. |
| 63 | // The maximum number of bytes that will be read is the size of the buffer. The minimum number |
| 64 | // of bytes that will be read is minBytes. If at least minBytes cannot be read, the promise |
| 65 | // will be resolved with the number of bytes read and the stream will be closed. |
| 66 | virtual kj::Promise<size_t> read(kj::ArrayPtr<kj::byte> buffer, size_t minBytes = 1) = 0; |
| 67 | |
| 68 | // If `end` is true, then `output.end()` will be called after pumping. Note that it's especially |
| 69 | // important to take advantage of this when using deferred proxying since calling `end()` |
| 70 | // directly might attempt to use the `IoContext` to call `registerPendingEvent()`. |
| 71 | // If the pump fails, the ReadableSource will be left in an errored state. The |
| 72 | // default implementation uses read() to read chunks of data and write them to the output |
| 73 | // using a 16KB buffer. |
| 74 | // |
| 75 | // Per the contract of pumpTo(), it is the caller's responsibility to ensure that both |
| 76 | // the WritableStreamSink and this ReadableSource remain alive until the returned |
| 77 | // promise resolves! |
| 78 | // |
| 79 | // It is the caller's responsibility to ensure that WritableStreamSink and this |
| 80 | // ReadableSource remain alive until the wrapped deferred proxy task resolves. |
| 81 | // The default implementation does arrange to make it safe to drop the source once |
| 82 | // the pump begins but that's only precautionary/defensive. It's still better/safer |
| 83 | // for the caller to keep the source alive until the pump completes. |
| 84 | virtual kj::Promise<DeferredProxy<void>> pumpTo( |
| 85 | WritableSink& output, EndAfterPump end = EndAfterPump::YES) = 0; |
| 86 | |
| 87 | // If the stream is still active, and the encoding matches an encoding that the stream |
| 88 | // can provide, gets the total length, if known. If the length is not known, or the |
| 89 | // encoding does not match the encoding of the underlying stream, or the stream is closed |
| 90 | // or errored, returns kj::none. |
| 91 | virtual kj::Maybe<size_t> tryGetLength(rpc::StreamEncoding encoding) = 0; |
| 92 | |
| 93 | // Fully consume the stream and return all of its data as a byte array. The limit |
| 94 | // parameter is the maximum number of bytes to read. If the stream contains more |
| 95 | // than this number of bytes, the promise will reject with an exception. |
| 96 | virtual kj::Promise<kj::Array<const kj::byte>> readAllBytes(size_t limit) = 0; |
| 97 | |
| 98 | // Fully consume the stream and return all of its data as a string. The limit |
| 99 | // parameter is the maximum number of bytes to read. If the stream contains more |
| 100 | // than this number of bytes, the promise will reject with an exception. |
| 101 | virtual kj::Promise<kj::String> readAllText(size_t limit) = 0; |
| 102 | |
| 103 | // Cancels the underlying source if it is still active. Must put the stream into an |
| 104 | // errored state. After calling this, all pending and future reads should fail. |
| 105 | // Dropping the ReadableSource without calling cancel() first should trigger |
| 106 | // cancel() with a generic exception. |
| 107 | virtual void cancel(kj::Exception reason) = 0; |
| 108 | |
| 109 | struct Tee { |
| 110 | kj::Own<ReadableSource> branch1; |
| 111 | kj::Own<ReadableSource> branch2; |
| 112 | }; |
| 113 | |
| 114 | // Tees the stream into two branches. The returned Tee contains two new ReadableSource |
| 115 | // instances that will each receive the same data. Once this is called, this instance is no |
| 116 | // longer usable and will behave as if it has been closed. |
| 117 | // The limit parameter specifies the maximum buffer size to use when teeing. |
| 118 | virtual Tee tee(size_t limit) = 0; |
| 119 | |
| 120 | // Gets the encoding of the stream. |
| 121 | virtual rpc::StreamEncoding getEncoding() = 0; |
| 122 | }; |
| 123 | |
| 124 | // Utility base class for ReadableSource wrappers that delegate all |
| 125 | // operations to an inner ReadableSource while selectively overriding |
| 126 | // some operations. |
| 127 | class ReadableSourceWrapper: public ReadableSource { |
| 128 | public: |
| 129 | KJ_DISALLOW_COPY_AND_MOVE(ReadableSourceWrapper); |
| 130 | virtual ~ReadableSourceWrapper() noexcept(false) = default; |
| 131 | |
| 132 | kj::Promise<size_t> read(kj::ArrayPtr<kj::byte> buffer, size_t minBytes = 1) override { |
| 133 | return getInner().read(buffer, minBytes); |
| 134 | } |
| 135 | |
| 136 | kj::Promise<DeferredProxy<void>> pumpTo( |
| 137 | WritableSink& output, EndAfterPump end = EndAfterPump::YES) override { |
| 138 | return getInner().pumpTo(output, end); |
| 139 | } |
| 140 | |
| 141 | kj::Promise<kj::Array<const kj::byte>> readAllBytes(size_t limit) override { |
| 142 | return getInner().readAllBytes(limit); |
| 143 | } |
| 144 | |
| 145 | kj::Promise<kj::String> readAllText(size_t limit) override { |
| 146 | return getInner().readAllText(limit); |
| 147 | } |
| 148 | |
| 149 | kj::Maybe<size_t> tryGetLength(rpc::StreamEncoding encoding) override { |
| 150 | return getInner().tryGetLength(encoding); |
| 151 | } |
| 152 | |
| 153 | void cancel(kj::Exception reason) override { |
| 154 | getInner().cancel(kj::mv(reason)); |
| 155 | } |
| 156 | |
| 157 | Tee tee(size_t limit) override { |
| 158 | return getInner().tee(limit); |
| 159 | } |
| 160 | |
| 161 | rpc::StreamEncoding getEncoding() override { |
| 162 | return getInner().getEncoding(); |
| 163 | } |
| 164 | |
| 165 | // Releases ownership of the inner ReadableSource. After calling this, |
| 166 | // this wrapper becomes unusable. |
| 167 | kj::Own<ReadableSource> release() { |
| 168 | auto ret = kj::mv(KJ_ASSERT_NONNULL(inner)); |
| 169 | inner = kj::none; |
| 170 | return kj::mv(ret); |
| 171 | } |
| 172 | |
| 173 | protected: |
| 174 | ReadableSourceWrapper(kj::Own<ReadableSource> inner): inner(kj::mv(inner)) {} |
| 175 | |
| 176 | ReadableSource& getInner() { |
| 177 | return *KJ_ASSERT_NONNULL(inner); |
| 178 | } |
| 179 | |
| 180 | private: |
| 181 | kj::Maybe<kj::Own<ReadableSource>> inner; |
| 182 | }; |
| 183 | |
| 184 | // Creates a ReadableSource that wraps the given kj::AsyncInputStream. |
| 185 | kj::Own<ReadableSource> newReadableSource(kj::Own<kj::AsyncInputStream> inner); |
| 186 | |
| 187 | // Creates a ReadableSource that is already in the errored state. |
| 188 | kj::Own<ReadableSource> newErroredReadableSource(kj::Exception exception); |
| 189 | |
| 190 | // Creates a ReadableSource that is already closed and will produce no data. |
| 191 | kj::Own<ReadableSource> newClosedReadableSource(); |
| 192 | |
| 193 | // Creates a ReadableSource that produces the given bytes and then closes. |
| 194 | // The backing object, if any, is held alive until the stream is closed or canceled. |
| 195 | // If the backing object is not provided, the bytes are copied. |
| 196 | kj::Own<ReadableSource> newReadableSourceFromBytes( |
| 197 | kj::ArrayPtr<const kj::byte> bytes, kj::Maybe<kj::Own<void>> backing = kj::none); |
| 198 | |
| 199 | // Creates a ReadableSource that wraps the given source and prevents deferred proxying. |
| 200 | kj::Own<ReadableSource> newIoContextWrappedReadableSource( |
| 201 | IoContext& ioctx, kj::Own<ReadableSource> inner); |
| 202 | |
| 203 | // Creates a ReadableSource that calls the given producer function to produce data |
| 204 | // on each read (useful primarily for testing). |
| 205 | kj::Own<ReadableSource> newReadableSourceFromProducer( |
| 206 | kj::Function<kj::Promise<size_t>(kj::ArrayPtr<kj::byte>, size_t)> producer, |
| 207 | kj::Maybe<uint64_t> expectedLength = kj::none); |
| 208 | |
| 209 | // Creates a ReadableSource that decodes the given stream according to the given encoding. |
| 210 | kj::Own<ReadableSource> newEncodedReadableSource( |
| 211 | rpc::StreamEncoding encoding, kj::Own<kj::AsyncInputStream> inner); |
| 212 | |
| 213 | // Wraps a kj::AsyncInputStream returned from a tee() call to ensure that it translates |
| 214 | // errors into equivalent JS exceptions. Typically this is used when customizing tee() on |
| 215 | // a ReadableSource implementation. |
| 216 | kj::Own<kj::AsyncInputStream> wrapTeeBranch(kj::Own<kj::AsyncInputStream> branch); |
| 217 | |
| 218 | // A ReadableStreamSource backed by in-memory data. Unlike newSystemStream() wrapping a |
| 219 | // newMemoryInputStream(), this implementation does NOT support deferred proxying. This is |
| 220 | // important when the backing memory has V8 heap provenance (e.g., jsg::BackingStore, Blob data, |
| 221 | // kj::Array<kj::byte> with a v8::BackingStore attached, etc) |
| 222 | // since the memory could be freed by GC after the IoContext completes. |
| 223 | // |
| 224 | // The `backing` parameter keeps the underlying memory alive for the lifetime of the stream. |
| 225 | // If not provided, the bytes are copied. |
| 226 | // |
| 227 | // TODO(soon): Update to implement ReadableSource instead of ReadableStreamSource. |
| 228 | // For now this is a ReadableStreamSource for compat with existing code. Once internal.h/c++ |
| 229 | // is updated to use ReadableSource, we will change this also. |
| 230 | // |
| 231 | // TODO(cleanup): It would be nice to eventually have some sort of stronger guarantee when |
| 232 | // deferred proxying can or cannot be used with a stream. Right now it's a bit ad hoc and |
| 233 | // error-prone. It requires the stream impl to keep track of whether it can be deferred-proxied |
| 234 | // or not, but in this case, that may be entirely opaque behind the details of the backing memory |
| 235 | // as is the case with kj::Array<kj::byte> instances that come from the type wrapper system. |
| 236 | kj::Own<ReadableStreamSource> newMemorySource( |
| 237 | kj::ArrayPtr<const kj::byte> bytes, kj::Maybe<kj::Own<void>> backing = kj::none); |
| 238 | |
| 239 | } // namespace api::streams |
| 240 | } // namespace workerd |