libutil: add async output stream type
we also extend AsyncInputStream with a drainInto variant to give async output streams rough feature parity with sync sinks. we still will not add serialization support to streams though, that's far too expensive. Change-Id: I60d5ab43610c45a40ea8740470a5eafe68064aea
This commit is contained in:
@@ -13,6 +13,18 @@ try {
|
||||
co_return result::current_exception();
|
||||
}
|
||||
|
||||
kj::Promise<Result<void>> AsyncInputStream::drainInto(AsyncOutputStream & stream)
|
||||
try {
|
||||
constexpr size_t BUF_SIZE = 65536;
|
||||
auto buf = std::make_unique<char[]>(BUF_SIZE);
|
||||
while (auto r = TRY_AWAIT(read(buf.get(), BUF_SIZE))) {
|
||||
TRY_AWAIT(stream.writeFull(buf.get(), r));
|
||||
}
|
||||
co_return result::success();
|
||||
} catch (...) {
|
||||
co_return result::current_exception();
|
||||
}
|
||||
|
||||
kj::Promise<Result<std::string>> AsyncInputStream::drain()
|
||||
try {
|
||||
StringSink s;
|
||||
|
||||
@@ -12,6 +12,8 @@
|
||||
|
||||
namespace nix {
|
||||
|
||||
class AsyncOutputStream;
|
||||
|
||||
// not derived from kj's AsyncInputStream because read and tryRead are already
|
||||
// taken as method names, we don't need the other functions, and the bit about
|
||||
// minBytes does not work well with our current io model. some day, who knows?
|
||||
@@ -24,6 +26,7 @@ public:
|
||||
virtual kj::Promise<Result<size_t>> read(void * buffer, size_t size) = 0;
|
||||
|
||||
kj::Promise<Result<void>> drainInto(Sink & sink);
|
||||
kj::Promise<Result<void>> drainInto(AsyncOutputStream & stream);
|
||||
|
||||
kj::Promise<Result<std::string>> drain();
|
||||
};
|
||||
@@ -77,4 +80,29 @@ public:
|
||||
|
||||
kj::Promise<Result<size_t>> read(void * data, size_t len) override;
|
||||
};
|
||||
|
||||
class AsyncOutputStream : private kj::AsyncObject
|
||||
{
|
||||
public:
|
||||
virtual ~AsyncOutputStream() noexcept(false) {}
|
||||
|
||||
virtual kj::Promise<Result<size_t>> write(const void * src, size_t size) = 0;
|
||||
|
||||
kj::Promise<Result<void>> writeFull(const void * src, size_t size)
|
||||
{
|
||||
return write(src, size).then(
|
||||
[this, src, size](Result<size_t> wrote) -> kj::Promise<Result<void>> {
|
||||
if (!wrote.has_value()) {
|
||||
return {wrote.error()};
|
||||
} else if (wrote.value() == size) {
|
||||
return {result::success()};
|
||||
} else {
|
||||
return writeFull(
|
||||
static_cast<const char *>(src) + wrote.value(), size - wrote.value()
|
||||
);
|
||||
}
|
||||
}
|
||||
);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user