libutil: add AsyncFramedInputStream
the async version of FramedSource, with all its weaknesses for bug-compat. Change-Id: Ifcdd4a5819f7cf25a0e8b01c63975ffead54079c
This commit is contained in:
@@ -835,7 +835,7 @@ kj::Promise<Result<void>> RemoteStore::ConnectionHandle::withFramedStream(
|
||||
)
|
||||
try {
|
||||
AsyncBufferedOutputStream to(stream);
|
||||
AsyncFramedStream sink(to);
|
||||
AsyncFramedOutputStream sink(to);
|
||||
|
||||
// NOLINTNEXTLINE(cppcoreguidelines-avoid-capturing-lambda-coroutines)
|
||||
auto send = [&]() -> kj::Promise<Result<void>> {
|
||||
|
||||
+59
-2
@@ -2,7 +2,9 @@
|
||||
#include "async.hh"
|
||||
#include "error.hh"
|
||||
#include "file-descriptor.hh"
|
||||
#include "logging.hh"
|
||||
#include "result.hh"
|
||||
#include "serialise.hh"
|
||||
#include <cerrno>
|
||||
#include <exception>
|
||||
#include <fcntl.h>
|
||||
@@ -195,8 +197,63 @@ kj::Promise<Result<size_t>> AsyncFdIoStream::write(const void * src, size_t size
|
||||
return {result::failure(std::make_exception_ptr(SysError(errno, "write failed")))};
|
||||
}
|
||||
}
|
||||
AsyncFramedInputStream::~AsyncFramedInputStream()
|
||||
{
|
||||
if (!eof) {
|
||||
printError(
|
||||
"AsyncFramedInputStream wasn't read to finish! its connection is now probably broken."
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
kj::Promise<Result<void>> AsyncFramedStream::finish()
|
||||
kj::Promise<Result<void>> AsyncFramedInputStream::finish()
|
||||
try {
|
||||
if (!eof) {
|
||||
while (true) {
|
||||
auto n = TRY_AWAIT(readNum<unsigned>(from));
|
||||
if (!n) {
|
||||
eof = true;
|
||||
break;
|
||||
}
|
||||
std::vector<char> data(n);
|
||||
if (TRY_AWAIT(from.readRange(data.data(), n, n)) != n) {
|
||||
throw Error("framed stream ended unexpectedly");
|
||||
}
|
||||
}
|
||||
}
|
||||
co_return result::success();
|
||||
} catch (...) {
|
||||
co_return result::current_exception();
|
||||
}
|
||||
|
||||
kj::Promise<Result<std::optional<size_t>>> AsyncFramedInputStream::read(void * buffer, size_t size)
|
||||
try {
|
||||
if (eof) {
|
||||
co_return std::nullopt;
|
||||
}
|
||||
|
||||
if (pos >= pending.size()) {
|
||||
size_t len = TRY_AWAIT(readNum<unsigned>(from));
|
||||
if (!len) {
|
||||
eof = true;
|
||||
co_return std::nullopt;
|
||||
}
|
||||
pending = std::vector<char>(len);
|
||||
pos = 0;
|
||||
if (TRY_AWAIT(from.readRange(pending.data(), len, len)) != len) {
|
||||
throw Error("framed stream ended unexpectedly");
|
||||
}
|
||||
}
|
||||
|
||||
auto n = std::min(size, pending.size() - pos);
|
||||
memcpy(buffer, pending.data() + pos, n);
|
||||
pos += n;
|
||||
co_return n;
|
||||
} catch (...) {
|
||||
co_return result::current_exception();
|
||||
}
|
||||
|
||||
kj::Promise<Result<void>> AsyncFramedOutputStream::finish()
|
||||
try {
|
||||
StringSink tmp;
|
||||
tmp << 0;
|
||||
@@ -206,7 +263,7 @@ try {
|
||||
co_return result::current_exception();
|
||||
}
|
||||
|
||||
kj::Promise<Result<size_t>> AsyncFramedStream::write(const void * buffer, size_t size)
|
||||
kj::Promise<Result<size_t>> AsyncFramedOutputStream::write(const void * buffer, size_t size)
|
||||
try {
|
||||
StringSink tmp;
|
||||
tmp << size;
|
||||
|
||||
+30
-2
@@ -207,15 +207,43 @@ public:
|
||||
kj::Promise<Result<size_t>> write(const void * src, size_t size) override;
|
||||
};
|
||||
|
||||
/**
|
||||
* A stream that reads a distinct format of concatenated chunks back into its
|
||||
* logical form, in order to guarantee a known state to the original stream,
|
||||
* even in the event of errors.
|
||||
*
|
||||
* Use with AsyncFramedOutputStream, which also allows the logical stream to be terminated
|
||||
* in the event of an exception.
|
||||
*/
|
||||
class AsyncFramedInputStream : public AsyncInputStream
|
||||
{
|
||||
private:
|
||||
AsyncInputStream & from;
|
||||
bool eof = false;
|
||||
/** Full contents of the current data frame. */
|
||||
std::vector<char> pending;
|
||||
/** Read offset into `pending`. The frame is fully processed if `pos == pending.size()`. */
|
||||
size_t pos = 0;
|
||||
|
||||
public:
|
||||
AsyncFramedInputStream(AsyncInputStream & from) : from(from) {}
|
||||
|
||||
~AsyncFramedInputStream();
|
||||
|
||||
kj::Promise<Result<void>> finish();
|
||||
|
||||
kj::Promise<Result<std::optional<size_t>>> read(void * buffer, size_t size) override;
|
||||
};
|
||||
|
||||
/**
|
||||
* Write as chunks in the format expected by FramedSource.
|
||||
*/
|
||||
class AsyncFramedStream : public AsyncOutputStream
|
||||
class AsyncFramedOutputStream : public AsyncOutputStream
|
||||
{
|
||||
AsyncOutputStream & to;
|
||||
|
||||
public:
|
||||
explicit AsyncFramedStream(AsyncOutputStream & to) : to(to) {}
|
||||
explicit AsyncFramedOutputStream(AsyncOutputStream & to) : to(to) {}
|
||||
|
||||
kj::Promise<Result<void>> finish();
|
||||
|
||||
|
||||
Reference in New Issue
Block a user