File
Blob: src/workerd/api/streams/writable-sink-adapter.h
| 1 | #include "common.h" |
| 2 | #include "writable-sink.h" |
| 3 | |
| 4 | #include <workerd/util/state-machine.h> |
| 5 | #include <workerd/util/weak-refs.h> |
| 6 | |
| 7 | namespace workerd::api::streams { |
| 8 | // Wraps a WritableStreamSink with a more JS-friendly interface that implements |
| 9 | // queued writes and backpressure signaling. This is arguably what WritableStreamSink |
| 10 | // should have been in the first place. Eventually we might be able to replace |
| 11 | // WritableStreamSink with this class directly, but for now we need to keep both. |
| 12 | // |
| 13 | // Instances of WritableStreamSinkJsAdapter are meant to be used from within the |
| 14 | // isolate lock, when you have need to write data to a kj stream from JavaScript. |
| 15 | // As such, it is not a jsg::Object itself, nor is it a kj I/O object, but it |
| 16 | // sits between the two worlds. Internally it holds the WritableStreamSink within |
| 17 | // an IoOwn so that correct IoContext usage is enforced. But the kj::Own for the |
| 18 | // adapter itself is meant to be held in JS land. |
| 19 | // |
| 20 | // Once created, the adapter owns the underlying WritableStreamSink. It is not |
| 21 | // possible to extract the sink from the adapter. This is because the adapter |
| 22 | // needs to be able to enforce its own state machine and queued write mechanism. |
| 23 | // |
| 24 | // The adapter implements backpressure signaling based on a high water mark |
| 25 | // configured at construction time. When the number of bytes in flight exceeds |
| 26 | // the high water mark, we signal backpressure by causing the ready promise |
| 27 | // to be reset to a new pending promise. When backpressure is released again, |
| 28 | // the ready promise is resolved. The identity of the ready promise changes |
| 29 | // whenever the backpressure state changes. |
| 30 | // |
| 31 | // The adapter also implements flush signaling. Flushing signals are checkpoints |
| 32 | // that are inserted into the write queue, essentially like a no-op write. They |
| 33 | // can be used as synchronization points to ensure that all prior writes have |
| 34 | // completed. Flush signals do not affect backpressure or stream state. |
| 35 | // |
| 36 | // Dropping the adapter will cancel any in-flight and pending operations |
| 37 | // immediately. Dropping the IoContext while the adapter is still active |
| 38 | // will also cancel any in-flight and pending operations and cause the |
| 39 | // adapter to be invalidated (the Active state is held with an IoOwn). |
| 40 | // |
| 41 | // ┌───────────────────────────────────────────┐ |
| 42 | // │ JavaScript Code │ |
| 43 | // │ │ |
| 44 | // │ • write(data) → Promise<void> │ |
| 45 | // │ • flush() → Promise<void> │ |
| 46 | // │ • end() → Promise<void> │ |
| 47 | // │ • abort(reason) │ |
| 48 | // │ • getReady() → Promise<void> │ |
| 49 | // └───────────────────────────────────────────┘ |
| 50 | // │ |
| 51 | // ▼ |
| 52 | // ┌───────────────────────────────────────────┐ |
| 53 | // │ WritableStreamSinkJsAdapter │ |
| 54 | // │ │ |
| 55 | // │ ┌─────────────────────────────────────┐ │ |
| 56 | // │ │ JavaScript API │ │ |
| 57 | // │ │ │ │ |
| 58 | // │ │ • write(data) → Promise<void> │ │ |
| 59 | // │ │ • flush() → Promise<void> │ │ |
| 60 | // │ │ • end() → Promise<void> │ │ |
| 61 | // │ │ • abort(reason) │ │ |
| 62 | // │ │ • getReady() → Promise<void> │ │ |
| 63 | // │ │ • getDesiredSize() → number │ │ |
| 64 | // │ └─────────────────────────────────────┘ │ |
| 65 | // │ │ │ |
| 66 | // │ ▼ │ |
| 67 | // │ ┌─────────────────────────────────────┐ │ |
| 68 | // │ │ Backpressure Management │ │ |
| 69 | // │ │ │ │ |
| 70 | // │ │ • High water mark (16KB default) │ │ |
| 71 | // │ │ • Bytes in flight tracking │ │ |
| 72 | // │ │ • Ready promise signaling │ │ |
| 73 | // │ │ • Queue depth management │ │ |
| 74 | // │ └─────────────────────────────────────┘ │ |
| 75 | // │ │ │ |
| 76 | // │ ▼ │ |
| 77 | // │ ┌─────────────────────────────────────┐ │ |
| 78 | // │ │ Write Queue Management │ │ |
| 79 | // │ │ │ │ |
| 80 | // │ │ • Queued writes with ordering │ │ |
| 81 | // │ │ • Flush checkpoints │ │ |
| 82 | // │ │ • Single in-flight write │ │ |
| 83 | // │ │ • Error propagation │ │ |
| 84 | // │ └─────────────────────────────────────┘ │ |
| 85 | // │ │ │ |
| 86 | // │ ▼ │ |
| 87 | // │ ┌─────────────────────────────────────┐ │ |
| 88 | // │ │ KJ Integration │ │ |
| 89 | // │ │ │ │ |
| 90 | // │ │ IoOwn<WritableStreamSink> │ │ |
| 91 | // │ │ WeakRef for safe references │ │ |
| 92 | // │ │ IoContext-aware operations │ │ |
| 93 | // │ └─────────────────────────────────────┘ │ |
| 94 | // └───────────────────────────────────────────┘ |
| 95 | // │ |
| 96 | // ▼ |
| 97 | // ┌───────────────────────────────────────────┐ |
| 98 | // │ WritableStreamSink │ |
| 99 | // │ (KJ Native Sink) │ |
| 100 | // │ │ |
| 101 | // │ • write(buffer) → Promise<void> │ |
| 102 | // │ • end() → Promise<void> │ |
| 103 | // │ • abort(reason) │ |
| 104 | // └───────────────────────────────────────────┘ |
| 105 | // |
| 106 | class WritableStreamSinkJsAdapter final { |
| 107 | public: |
| 108 | struct Options { |
| 109 | // While the WritableStreamSink interface, and kj streams in general, do |
| 110 | // not have a notion of backpressure, and instead generally require only |
| 111 | // one write to be in flight at a time, it's better for performance for |
| 112 | // us to be able to buffer a bit more data in flight. So we will implement |
| 113 | // a simple high water mark mechanism. The default is 16KB. |
| 114 | size_t highWaterMark = 16384; |
| 115 | |
| 116 | // When detachOnWrite is true, and a write() is made with an ArrayBuffer, |
| 117 | // or ArrayBufferView, we will attempt to detach the underlying buffer |
| 118 | // before writing it to the sink. Detaching is required by the |
| 119 | // streams spec but our original implementation does not detach |
| 120 | // and it turns out there are old workers depending on that behavior. |
| 121 | bool detachOnWrite = false; |
| 122 | }; |
| 123 | |
| 124 | WritableStreamSinkJsAdapter(jsg::Lock& js, |
| 125 | IoContext& ioContext, |
| 126 | kj::Own<WritableSink> sink, |
| 127 | kj::Maybe<Options> options = kj::none); |
| 128 | WritableStreamSinkJsAdapter(jsg::Lock& js, |
| 129 | IoContext& ioContext, |
| 130 | kj::Own<kj::AsyncOutputStream> stream, |
| 131 | StreamEncoding encoding, |
| 132 | kj::Maybe<Options> options = kj::none); |
| 133 | KJ_DISALLOW_COPY_AND_MOVE(WritableStreamSinkJsAdapter); |
| 134 | ~WritableStreamSinkJsAdapter() noexcept(false); |
| 135 | |
| 136 | // If we are in the errored state, returns the exception, otherwise kj::none. |
| 137 | kj::Maybe<const kj::Exception&> isErrored() KJ_LIFETIMEBOUND; |
| 138 | |
| 139 | // Returns true if we are in the closed state. |
| 140 | bool isClosed(); |
| 141 | |
| 142 | // Returns true if close() has been called but we are not yet closed. |
| 143 | bool isClosing(); |
| 144 | |
| 145 | // If we are not in the closed or errored state, returns the desired |
| 146 | // size based on the configured high water mark and the number of |
| 147 | // bytes currently in flight. The desired size is the number of bytes |
| 148 | // that can be written before we exceed the high water mark. If the |
| 149 | // return value is <= 0 then backpressure is being signaled. If we are |
| 150 | // in the closed or errored states, returns kj::none. |
| 151 | kj::Maybe<ssize_t> getDesiredSize(); |
| 152 | |
| 153 | // Writes a chunk to the underlying sink via the queued write mechanism. |
| 154 | // The implementation ensures that only one write is in flight with the |
| 155 | // underlying sink at a time, while additional writes are queued up |
| 156 | // behind it. It is not necessary to await the returned promise before |
| 157 | // calling write() again, though doing so is not an error. If the write |
| 158 | // fails, the returned promise will reject with the failure reason. |
| 159 | // Also if the write fails, the adapter will be transitioned to the |
| 160 | // errored state and all subsequent queued writes will fail. Once |
| 161 | // close() has been called, no additional writes will be accepted |
| 162 | // and the returned promise will reject with an error. If the adapter |
| 163 | // is already in the closed or errored state, the returned promise will |
| 164 | // be rejected. |
| 165 | // |
| 166 | // Values written may be ArrayBuffer, ArrayBufferView, SharedArrayBuffer, |
| 167 | // or string. Other types will cause the returned promise to reject. |
| 168 | // |
| 169 | // Backpressure is signaled when the number of bytes in flight (i.e. |
| 170 | // the total number of bytes passed to write() calls that have not yet |
| 171 | // completed) exceeds the configured high water mark. When backpressure |
| 172 | // is signaled, additional writes are still accepted and queued up, but |
| 173 | // the caller really should wait for the ready promise to resolve before |
| 174 | // continuing to write more. This works exactly like a WritableStream's |
| 175 | // backpressure mechanism. Callers keep writing until backpressure is |
| 176 | // signaled, then wait for the ready promise to resolve before continuing, |
| 177 | // etc. |
| 178 | jsg::Promise<void> write(jsg::Lock& js, const jsg::JsValue& value); |
| 179 | |
| 180 | // Inserts a flush signal into the write queue. The returned promise |
| 181 | // resolves once all prior writes have completed. This can be used |
| 182 | // as a synchronization point to ensure that all writes up to this |
| 183 | // point have been fully processed. If the adapter is in the closed |
| 184 | // or errored state, the returned promise will reject. If the stream |
| 185 | // errors while waiting for prior writes to complete, the returned |
| 186 | // promise will be rejected. |
| 187 | jsg::Promise<void> flush(jsg::Lock& js); |
| 188 | |
| 189 | // Transitions the adapter into the closing state. Once the write queue |
| 190 | // is empty, we will close the sink and transition to the closed state. |
| 191 | // If the adapter is already in the closing state, a new promise is |
| 192 | // returned that will resolve when the adapter is fully closed. If the |
| 193 | // adapter is already closed, a resolved promise is returned. If the |
| 194 | // adapter is in the errored state, a rejected promise is returned. |
| 195 | // All pending writes in the queue will be processed before closing |
| 196 | // the sink and transitioning to the closed state. If any pending |
| 197 | // writes fail, the adapter will transition to the errored state, and |
| 198 | // all subsequent pending writes will be rejected along with the close |
| 199 | // promise. |
| 200 | jsg::Promise<void> end(jsg::Lock& js); |
| 201 | |
| 202 | // Transitions the adapter to the errored state, even if we are already closed. |
| 203 | // All pending or in-flight writes, and a pending close, will all be rejected |
| 204 | // with the given exception. If we are already in the errored state, this |
| 205 | // is a no-op. This change is immediate. Once in the errored state, no |
| 206 | // further writes or closes are allowed. |
| 207 | void abort(kj::Exception&& exception); |
| 208 | |
| 209 | // Transitions the adapter to the errored state, even if we are already closed. |
| 210 | // All pending or in-flight writes, and a pending close, will all be rejected |
| 211 | // with the given exception. If we are already in the errored state, this |
| 212 | // is a no-op. This change is immediate. Once in the errored state, no |
| 213 | // further writes or closes are allowed. This variant is for use when |
| 214 | // the exception is coming from JavaScript. It will be converted into a |
| 215 | // tunneled kj::Exception. |
| 216 | void abort(jsg::Lock& js, const jsg::JsValue& reason); |
| 217 | |
| 218 | // Returns a promise that resolves when backpressure is released. |
| 219 | // Note that the identity of the returned promise will change as the |
| 220 | // backpressure state changes. Whenever backpressure is signaled, a new |
| 221 | // pending promise will be created, whenever backpressure is released |
| 222 | // again that promise will be resolved. As such, this promise should |
| 223 | // not be cached or stored. Instead, before every write() call, the |
| 224 | // caller should wait on the current getReady() promise. |
| 225 | jsg::Promise<void> getReady(jsg::Lock& js); |
| 226 | |
| 227 | // Returns a memoized identity for the ready promise. This can be used |
| 228 | // to return a stable reference to the ready promise out to JavaScript |
| 229 | // that will not change identity between calls unless the backpressure |
| 230 | // state changes. Like the getReady() promise, this should not be cached |
| 231 | // or stored, but it is safe to return this from a getter multiple times |
| 232 | // to JavaScript as it will ensure that the same JS promise object is |
| 233 | // always returned until the backpressure state changes. This variation |
| 234 | // is not suitable for use within C++ code that needs to await on the |
| 235 | // ready promise because the internal jsg::Promise<void> object will |
| 236 | // no longer exist once the reference is passed out to JavaScript. |
| 237 | jsg::MemoizedIdentity<jsg::Promise<void>>& getReadyStable(); |
| 238 | |
| 239 | // Returns the options used to configure this adapter if the adapter |
| 240 | // is not closed or errored. |
| 241 | kj::Maybe<const Options&> getOptions(); |
| 242 | |
| 243 | void visitForGc(jsg::GcVisitor& visitor); |
| 244 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const; |
| 245 | |
| 246 | private: |
| 247 | // Represents the active state of the adapter. Importantly, this state |
| 248 | // holds both the underlying WritableStreamSink and the write queue. |
| 249 | // It must be held within an IoOwn. |
| 250 | struct Active; |
| 251 | |
| 252 | struct Closed final { |
| 253 | static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj; |
| 254 | }; |
| 255 | |
| 256 | struct Open { |
| 257 | static constexpr kj::StringPtr NAME KJ_UNUSED = "open"_kj; |
| 258 | IoOwn<Active> active; |
| 259 | }; |
| 260 | |
| 261 | // State machine for tracking writable sink adapter lifecycle: |
| 262 | // Open -> Closed (normal close via end()) |
| 263 | // Open -> kj::Exception (error via abort() or write failure) |
| 264 | // Closed is terminal, kj::Exception is implicitly terminal via ErrorState. |
| 265 | using State = StateMachine<TerminalStates<Closed>, |
| 266 | ErrorState<kj::Exception>, |
| 267 | ActiveState<Open>, |
| 268 | Open, |
| 269 | Closed, |
| 270 | kj::Exception>; |
| 271 | State state; |
| 272 | |
| 273 | // Used for backpressure signaling. When backpressure is indicated, the |
| 274 | // readyResolver, ready, and readyWatcher will be replaced with a new set. |
| 275 | // When backpressure is relieved, the readyResolver will be resolved. |
| 276 | // The adapter will start out in a ready state. |
| 277 | struct BackpressureState final { |
| 278 | // Note that if the BackpressureState is dropped while in a waiting state, |
| 279 | // the ready promise will be left unresolved. This is OK. |
| 280 | kj::Maybe<jsg::Promise<void>::Resolver> readyResolver; |
| 281 | jsg::Promise<void> ready; |
| 282 | jsg::MemoizedIdentity<jsg::Promise<void>> readyWatcher; |
| 283 | |
| 284 | // Aborts backpressure signaling, likely because the adapter is being errored. |
| 285 | // Causes the ready promise to be rejected with the given reason. |
| 286 | void abort(jsg::Lock& js, const jsg::JsValue& reason); |
| 287 | |
| 288 | // Releases backpressure, resolving the ready promise. |
| 289 | void release(jsg::Lock& js); |
| 290 | |
| 291 | // Indicates that backpressure has been signaled and we are waiting |
| 292 | // for it to be released or aborted. |
| 293 | bool isWaiting() const; |
| 294 | |
| 295 | // Returns a promise that resolves when backpressure is released. |
| 296 | // Note that every call to this returns a new jsg::Promise<void> |
| 297 | // instance. Callers that need a stable identity should use |
| 298 | // getReadyStable() instead (generally this is only the case when |
| 299 | // returning the promise to JavaScript via a getter). |
| 300 | jsg::Promise<void> getReady(jsg::Lock& js); |
| 301 | |
| 302 | // Returns a memoized identity for the ready promise. This can be used |
| 303 | // to return a stable reference to the ready promise out to JavaScript |
| 304 | // that will not change identity between calls unless the backpressure |
| 305 | // state changes. |
| 306 | jsg::MemoizedIdentity<jsg::Promise<void>>& getReadyStable(); |
| 307 | BackpressureState(jsg::Promise<void>::Resolver&& resolver, |
| 308 | jsg::Promise<void>&& promise, |
| 309 | jsg::MemoizedIdentity<jsg::Promise<void>>&& watcher); |
| 310 | }; |
| 311 | BackpressureState backpressureState; |
| 312 | kj::Rc<WeakRef<WritableStreamSinkJsAdapter>> selfRef; |
| 313 | |
| 314 | // Replaces the backpressure state with a new one, indicating that backpressure |
| 315 | // is being applied. If we are already in a backpressure state, this is a no-op. |
| 316 | // This will cause the ready promise (and its stable identity) to change. |
| 317 | void maybeSignalBackpressure(jsg::Lock& js); |
| 318 | |
| 319 | // Conditionally releases backpressure if the desired size is now > 0. |
| 320 | void maybeReleaseBackpressure(jsg::Lock& js); |
| 321 | |
| 322 | // Creates a new BackpressureState in the waiting state. |
| 323 | static BackpressureState newBackpressureState(jsg::Lock& js); |
| 324 | }; |
| 325 | |
| 326 | // ================================================================================ |
| 327 | |
| 328 | // Adapts a WritableStream to a KJ-frendly interface. |
| 329 | // The adapter fully wraps the WritableStream instance, |
| 330 | // using a WritableStreamDefaultWriter to push data to it. |
| 331 | // Then the adapter is destroyed or aborted, the writer is |
| 332 | // aborted and both the writer and the stream references |
| 333 | // are dropped. Critically, the stream is not usable after |
| 334 | // ownership is transferred to this adapter. Initializing the adapter |
| 335 | // will fail if the stream is already locked. |
| 336 | // |
| 337 | // If the adapter is dropped, or aborted while there are pending writes, |
| 338 | // the pending writes will be rejected with the same exception as the abort. |
| 339 | // |
| 340 | // While WritableStream itself allows multiple writes to be in flight |
| 341 | // at the same time, the WritableStreamSink interface does not, so |
| 342 | // the adapter will ensure that only one write is in flight at a time. |
| 343 | // |
| 344 | // While the caller is expected to follow the WritableStreamSink contract |
| 345 | // and keep the adapter alive until the write promises resolve, there |
| 346 | // are some protections in place to avoid use-after-free if the caller |
| 347 | // drops the adapter. There's nothing we can do if the caller drops the |
| 348 | // buffer, however, so that is still a hard requirement. |
| 349 | // TODO(safety): This can be made safer by having write take a kj::Array |
| 350 | // as input instead of a kj::ArrayPtr but that's a larger refactor. |
| 351 | // |
| 352 | // ┌───────────────────────────────────────────┐ |
| 353 | // │ WritableStreamSink │ |
| 354 | // │ │ |
| 355 | // │ • write(buffer) │ |
| 356 | // │ • write(pieces[]) │ |
| 357 | // │ • end() │ |
| 358 | // │ • abort(reason) │ |
| 359 | // └───────────────────────────────────────────┘ |
| 360 | // │ |
| 361 | // ▼ |
| 362 | // ┌───────────────────────────────────────────┐ |
| 363 | // │ WritableStreamSinkKjAdapter │ |
| 364 | // │ │ |
| 365 | // │ ┌─────────────────────────────────────┐ │ |
| 366 | // │ │ KJ Native API │ │ |
| 367 | // │ │ │ │ |
| 368 | // │ │ • write(ArrayPtr<byte>) │ │ |
| 369 | // │ │ • write(ArrayPtr<ArrayPtr<byte>>) │ │ |
| 370 | // │ │ • end() → Promise<void> │ │ |
| 371 | // │ │ • abort(exception) │ │ |
| 372 | // │ └─────────────────────────────────────┘ │ |
| 373 | // │ │ │ |
| 374 | // │ ▼ │ |
| 375 | // │ ┌─────────────────────────────────────┐ │ |
| 376 | // │ │ State Management │ │ |
| 377 | // │ │ │ │ |
| 378 | // │ │ Active ──► Closed │ │ |
| 379 | // │ │ │ │ │ │ |
| 380 | // │ │ │ ▼ │ │ |
| 381 | // │ │ └─────► Errored │ │ |
| 382 | // │ └─────────────────────────────────────┘ │ |
| 383 | // │ │ │ |
| 384 | // │ ▼ │ |
| 385 | // │ ┌─────────────────────────────────────┐ │ |
| 386 | // │ │ JavaScript Integration │ │ |
| 387 | // │ │ │ │ |
| 388 | // │ │ WritableStreamDefaultWriter │ │ |
| 389 | // │ │ WeakRef for safe references │ │ |
| 390 | // │ │ IoContext-aware JS operations │ │ |
| 391 | // │ │ Promise handling & async writes │ │ |
| 392 | // │ └─────────────────────────────────────┘ │ |
| 393 | // └───────────────────────────────────────────┘ |
| 394 | // │ |
| 395 | // ▼ |
| 396 | // ┌───────────────────────────────────────────┐ |
| 397 | // │ JavaScript WritableStream │ |
| 398 | // │ │ |
| 399 | // │ • getWriter() │ |
| 400 | // │ • write(chunk) → Promise<void> │ |
| 401 | // │ • close() → Promise<void> │ |
| 402 | // │ • abort(reason) → Promise<void> │ |
| 403 | // │ • locked, state properties │ |
| 404 | // └───────────────────────────────────────────┘ |
| 405 | // |
| 406 | class WritableStreamSinkKjAdapter final: public WritableSink { |
| 407 | public: |
| 408 | WritableStreamSinkKjAdapter(jsg::Lock& js, IoContext& ioContext, jsg::Ref<WritableStream> stream); |
| 409 | ~WritableStreamSinkKjAdapter() noexcept(false); |
| 410 | |
| 411 | // Attempts to write the given buffer to the underlying stream. |
| 412 | // The returned promise resolves once the write has completed. |
| 413 | // If the stream is closed, the returned promise rejects with |
| 414 | // an exception. If the stream errors, the returned promise |
| 415 | // rejects with the same exception. If the write fails, the |
| 416 | // returned promise rejects with the failure reason. |
| 417 | // |
| 418 | // Per the contract of write, it is the caller's responsibility |
| 419 | // to ensure that the adapter and buffer remain alive until |
| 420 | // the returned promise resolves. |
| 421 | kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override; |
| 422 | |
| 423 | // Attempts to write the given pieces to the underlying stream. |
| 424 | // The returned promise resolves once the full write has completed. |
| 425 | // If the stream is closed, the returned promise rejects with |
| 426 | // an exception. If the stream errors, the returned promise |
| 427 | // rejects with the same exception. If the write fails, the |
| 428 | // returned promise rejects with the failure reason. |
| 429 | // Per the contract of write, it is the caller's responsibility |
| 430 | // to ensure that the adapter and buffers remain alive until |
| 431 | // the returned promise resolves. |
| 432 | kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override; |
| 433 | |
| 434 | // Closes the underlying stream. The returned promise resolves |
| 435 | // once the stream is fully closed. If the stream is already |
| 436 | // closed, the returned promise resolves immediately. If the |
| 437 | // stream errors, the returned promise rejects with the same |
| 438 | // exception. If the close fails, the returned promise rejects |
| 439 | // with the failure reason. |
| 440 | kj::Promise<void> end() override; |
| 441 | |
| 442 | // Immediately interrupts existing pending writes and errors the stream. |
| 443 | // All pending or in-flight writes will be rejected with the given |
| 444 | // exception. If we are already in the errored state, this is a no-op |
| 445 | // and the exception is ignored. This change is immediate. Once in |
| 446 | // the errored state, no further writes or closes are allowed. |
| 447 | void abort(kj::Exception reason) override; |
| 448 | |
| 449 | // A WritableStreamSinkKjAdapter always (currently) uses identity encoding. |
| 450 | rpc::StreamEncoding disownEncodingResponsibility() override { |
| 451 | return rpc::StreamEncoding::IDENTITY; |
| 452 | } |
| 453 | |
| 454 | rpc::StreamEncoding getEncoding() override { |
| 455 | return rpc::StreamEncoding::IDENTITY; |
| 456 | } |
| 457 | |
| 458 | private: |
| 459 | struct Active; |
| 460 | KJ_DECLARE_NON_POLYMORPHIC(Active); |
| 461 | |
| 462 | struct KjClosed { |
| 463 | static constexpr kj::StringPtr NAME KJ_UNUSED = "closed"_kj; |
| 464 | }; |
| 465 | |
| 466 | struct KjOpen { |
| 467 | static constexpr kj::StringPtr NAME KJ_UNUSED = "open"_kj; |
| 468 | kj::Own<Active> active; |
| 469 | }; |
| 470 | |
| 471 | // State machine for tracking writable sink adapter lifecycle: |
| 472 | // KjOpen -> KjClosed (normal close via end()) |
| 473 | // KjOpen -> kj::Exception (error via abort() or write failure) |
| 474 | // KjClosed is terminal, kj::Exception is implicitly terminal via ErrorState. |
| 475 | using KjState = StateMachine<TerminalStates<KjClosed>, |
| 476 | ErrorState<kj::Exception>, |
| 477 | ActiveState<KjOpen>, |
| 478 | KjOpen, |
| 479 | KjClosed, |
| 480 | kj::Exception>; |
| 481 | KjState state; |
| 482 | kj::Rc<WeakRef<WritableStreamSinkKjAdapter>> selfRef; |
| 483 | }; |
| 484 | |
| 485 | } // namespace workerd::api::streams |