libutil: add buffered async streams

these will let us share async stream io buffers with sync sinks and sources.

Change-Id: If3149803a9e1fda62391399177da62f7522a811b
This commit is contained in:
eldritch horrors
2025-06-17 14:34:05 +02:00
parent 81d2d26c3f
commit 7b65d7c508
2 changed files with 97 additions and 0 deletions
+50
View File
@@ -83,4 +83,54 @@ try {
} catch (...) {
return {result::current_exception()};
}
kj::Promise<Result<size_t>> AsyncBufferedInputStream::read(void * data, size_t size)
try {
while (buffer->used() == 0) {
const auto space = buffer->getWriteBuffer();
const auto got = TRY_AWAIT(inner.read(space.data(), space.size()));
if (got == 0) {
co_return 0;
}
buffer->added(got);
}
const auto available = buffer->getReadBuffer();
const auto n = std::min(size, available.size());
memcpy(data, available.data(), n);
buffer->consumed(n);
co_return result::success(n);
} catch (...) {
co_return result::current_exception();
}
kj::Promise<Result<size_t>> AsyncBufferedOutputStream::write(const void * src, size_t size)
try {
if (size > buffer->size()) {
TRY_AWAIT(flush());
co_return TRY_AWAIT(inner.write(src, size));
}
if (size > buffer->getWriteBuffer().size()) {
TRY_AWAIT(flush());
}
const auto into = buffer->getWriteBuffer();
memcpy(into.data(), src, size);
buffer->added(size);
co_return size;
} catch (...) {
co_return result::current_exception();
}
kj::Promise<Result<void>> AsyncBufferedOutputStream::flush()
try {
if (auto unsent = buffer->getReadBuffer(); !unsent.empty()) {
TRY_AWAIT(inner.writeFull(unsent.data(), unsent.size()));
buffer->consumed(unsent.size());
}
co_return result::success();
} catch (...) {
co_return result::current_exception();
}
}
+47
View File
@@ -3,6 +3,8 @@
#include "lix/libutil/async.hh"
#include "lix/libutil/box_ptr.hh"
#include "lix/libutil/io-buffer.hh"
#include "lix/libutil/ref.hh"
#include "lix/libutil/result.hh"
#include "lix/libutil/serialise.hh"
#include <kj/async.h>
@@ -81,6 +83,28 @@ public:
kj::Promise<Result<size_t>> read(void * data, size_t len) override;
};
class AsyncBufferedInputStream : public AsyncInputStream
{
AsyncInputStream & inner;
ref<IoBuffer> buffer;
public:
AsyncBufferedInputStream(AsyncInputStream & inner, ref<IoBuffer> buffer)
: inner(inner)
, buffer(buffer)
{
}
AsyncBufferedInputStream(AsyncInputStream & inner, size_t bufSize = 32 * 1024)
: AsyncBufferedInputStream(inner, make_ref<IoBuffer>(bufSize))
{
}
KJ_DISALLOW_COPY_AND_MOVE(AsyncBufferedInputStream);
kj::Promise<Result<size_t>> read(void * data, size_t size) override;
};
class AsyncOutputStream : private kj::AsyncObject
{
public:
@@ -105,4 +129,27 @@ public:
);
}
};
class AsyncBufferedOutputStream : public AsyncOutputStream
{
AsyncOutputStream & inner;
ref<IoBuffer> buffer;
public:
AsyncBufferedOutputStream(AsyncOutputStream & inner, ref<IoBuffer> buffer)
: inner(inner)
, buffer(buffer)
{
}
AsyncBufferedOutputStream(AsyncOutputStream & inner, size_t bufSize = 32 * 1024)
: AsyncBufferedOutputStream(inner, make_ref<IoBuffer>(bufSize))
{
}
KJ_DISALLOW_COPY_AND_MOVE(AsyncBufferedOutputStream);
kj::Promise<Result<size_t>> write(const void * src, size_t size) override;
kj::Promise<Result<void>> flush();
};
}