File
Blob: src/workerd/util/stream-utils.h
| 1 | #pragma once |
| 2 | |
| 3 | #include <kj/async-io.h> |
| 4 | |
| 5 | namespace workerd { |
| 6 | |
| 7 | kj::Own<kj::AsyncIoStream> newNullIoStream(); |
| 8 | kj::Own<kj::AsyncInputStream> newNullInputStream(); |
| 9 | kj::Own<kj::AsyncOutputStream> newNullOutputStream(); |
| 10 | |
| 11 | // Get a shared global null output stream (singleton, thread-safe) |
| 12 | kj::AsyncOutputStream& getGlobalNullOutputStream(); |
| 13 | |
| 14 | // When maybeBacking is provided, it is held onto by the MemoryInputStream |
| 15 | // using a kj::Rc<...> so that teeing the stream can share ownership of |
| 16 | // the backing storage without the need for any additional buffering. If |
| 17 | // the backing storage is not provided, then optimized teeing of the stream |
| 18 | // will not be supported and the implementation will return a kj::none from |
| 19 | // tryTee(). |
| 20 | kj::Own<kj::AsyncInputStream> newMemoryInputStream( |
| 21 | kj::ArrayPtr<const kj::byte>, kj::Maybe<kj::Own<void>> maybeBacking = kj::none); |
| 22 | kj::Own<kj::AsyncInputStream> newMemoryInputStream( |
| 23 | kj::StringPtr, kj::Maybe<kj::Own<void>> maybeBacking = kj::none); |
| 24 | |
| 25 | // An InputStream that can be disconnected. |
| 26 | class NeuterableInputStream: public kj::AsyncInputStream, public kj::Refcounted { |
| 27 | public: |
| 28 | virtual void neuter(kj::Exception ex) = 0; |
| 29 | }; |
| 30 | |
| 31 | class NeuterableIoStream: public kj::AsyncIoStream { |
| 32 | public: |
| 33 | virtual void neuter(kj::Exception ex) = 0; |
| 34 | }; |
| 35 | |
| 36 | // Until kj::AsyncOutputStream has an end() method of its own... We |
| 37 | // provide this subclass that adds it. |
| 38 | class EndableAsyncOutputStream: public kj::AsyncOutputStream { |
| 39 | public: |
| 40 | // By default, end() is a no-op. Subclasses may override. |
| 41 | virtual kj::Promise<void> end() { |
| 42 | co_return; |
| 43 | } |
| 44 | }; |
| 45 | |
| 46 | kj::Own<NeuterableInputStream> newNeuterableInputStream(kj::AsyncInputStream&); |
| 47 | kj::Own<NeuterableIoStream> newNeuterableIoStream(kj::AsyncIoStream&); |
| 48 | |
| 49 | } // namespace workerd |