Skip to content
File

Blob: src/workerd/api/streams/writable-sink.h

cpp143 lines
1#pragma once
2 
3#include <workerd/io/worker-interface.capnp.h>
4 
5#include <kj/debug.h>
6 
7namespace kj {
8class AsyncOutputStream;
9}
10 
11namespace workerd {
12 
13class IoContext;
14 
15namespace api::streams {
16 
17// A WritableSink is primarily intended to serve as a bridge between kj::AsyncOutputStream
18// and the WritableStream API. However, it can also be used directly by KJ-space code. While
19// WritableSink should probably have been a more JS-friendly API, it's a bit too late
20// to change that now. Use the WritableSinkJsAdapter in the writable-sink-adapter.h file
21// to wrap a WritableSink for use from JavaScript.
22//
23// Not all WritableSink implementations will be explicitly backed by a KJ stream;
24// some might be test implementations that discard data or accumulate it in memory, for
25// instance.
26//
27// A WritableSink must be treated like a KJ I/O object. Instances that are held
28// by any JS-heap objects must be held by an IoOwn.
29//
30// The sink permits only one write() or end() operation to be pending at a time. If
31// a second write() or end() is attempted while one is already pending, the promise
32// returned by the second call will be rejected with a jsg::Error. This is to
33// match the behavior of the kj::AsyncOutputStream interface.
34//
35// If the sink is aborted or dropped, any pending write() or end() operations will be
36// canceled.
37class WritableSink {
38 public:
39 // Write the given buffer to the stream, returning a promise that resolves when the write
40 // completes.
41 virtual kj::Promise<void> write(kj::ArrayPtr<const kj::byte> buffer) KJ_WARN_UNUSED_RESULT = 0;
42 
43 // Write the given pieces to the stream, returning a promise that resolves when the write
44 // completes.
45 virtual kj::Promise<void> write(
46 kj::ArrayPtr<const kj::ArrayPtr<const kj::byte>> pieces) KJ_WARN_UNUSED_RESULT = 0;
47 
48 // Ends the stream, transitioning it to the closed state. After this, no further writes
49 // will be accepted.
50 virtual kj::Promise<void> end() KJ_WARN_UNUSED_RESULT = 0;
51 
52 // Aborts the stream, transitioning it to the errored state. After this, no further writes
53 // will be accepted.
54 virtual void abort(kj::Exception reason) = 0;
55 
56 // Tells the sink that it is no longer to be responsible for encoding in the correct format.
57 // Instead, the caller takes responsibility. The expected encoding is returned; the caller
58 // promises that all future writes will use this encoding.
59 virtual rpc::StreamEncoding disownEncodingResponsibility() = 0;
60 
61 // Return the encoding that this sink is using.
62 virtual rpc::StreamEncoding getEncoding() = 0;
63};
64 
65// Utility base class for WritableSink wrappers that delegate all
66// operations to an inner WritableSink while selectively overriding
67// some operations.
68class WritableSinkWrapper: public WritableSink {
69 public:
70 virtual ~WritableSinkWrapper() noexcept(false) {
71 canceler.cancel(KJ_EXCEPTION(DISCONNECTED, "Dropped"));
72 }
73 
74 kj::Promise<void> write(kj::ArrayPtr<const kj::byte> buffer) override {
75 return getInner().write(buffer);
76 }
77 
78 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const kj::byte>> pieces) override {
79 return getInner().write(pieces);
80 }
81 
82 kj::Promise<void> end() override {
83 return getInner().end();
84 }
85 
86 void abort(kj::Exception reason) override {
87 getInner().abort(kj::mv(reason));
88 }
89 
90 rpc::StreamEncoding disownEncodingResponsibility() override {
91 return getInner().disownEncodingResponsibility();
92 }
93 
94 rpc::StreamEncoding getEncoding() override {
95 return getInner().getEncoding();
96 }
97 
98 // Releases ownership of the inner WritableSink. After calling this,
99 // this instance is no longer usable.
100 kj::Own<WritableSink> release() {
101 auto ret = kj::mv(KJ_ASSERT_NONNULL(inner));
102 inner = kj::none;
103 canceler.cancel(KJ_EXCEPTION(DISCONNECTED, "Released"));
104 return kj::mv(ret);
105 }
106 
107 protected:
108 WritableSinkWrapper(kj::Own<WritableSink> inner): inner(kj::mv(inner)) {}
109 KJ_DISALLOW_COPY_AND_MOVE(WritableSinkWrapper);
110 
111 WritableSink& getInner() {
112 return *KJ_ASSERT_NONNULL(inner);
113 }
114 
115 private:
116 kj::Canceler canceler;
117 kj::Maybe<kj::Own<WritableSink>> inner;
118};
119 
120// Creates a WritableSink that wraps a kj::AsyncOutputStream.
121kj::Own<WritableSink> newWritableSink(kj::Own<kj::AsyncOutputStream> inner);
122 
123// Creates a WritableSink that is in the closed state.
124kj::Own<WritableSink> newClosedWritableSink();
125 
126// Creates a WritableSink that is permanently in the errored state.
127kj::Own<WritableSink> newErroredWritableSink(kj::Exception reason);
128 
129// Creates a WritableSink that discards all data written to it.
130kj::Own<WritableSink> newNullWritableSink();
131 
132// Creates a WritableSink that encodes data written to it.
133kj::Own<WritableSink> newEncodedWritableSink(
134 rpc::StreamEncoding encoding, kj::Own<kj::AsyncOutputStream> inner);
135 
136// Wraps a WritableSink such that each write()/end() call on the returned sink will
137// register as a pending event on the IoContext.
138kj::Own<WritableSink> newIoContextWrappedWritableSink(
139 IoContext& ioContext, kj::Own<WritableSink> inner);
140 
141} // namespace api::streams
142} // namespace workerd