From 7b65d7c508cbe0f7f0ae971ab5c9ac0c7a7eb69e Mon Sep 17 00:00:00 2001 From: eldritch horrors Date: Mon, 16 Jun 2025 18:51:59 +0200 Subject: [PATCH] libutil: add buffered async streams these will let us share async stream io buffers with sync sinks and sources. Change-Id: If3149803a9e1fda62391399177da62f7522a811b --- lix/libutil/async-io.cc | 50 +++++++++++++++++++++++++++++++++++++++++ lix/libutil/async-io.hh | 47 ++++++++++++++++++++++++++++++++++++++ 2 files changed, 97 insertions(+) diff --git a/lix/libutil/async-io.cc b/lix/libutil/async-io.cc index eaa6786c6..22d943437 100644 --- a/lix/libutil/async-io.cc +++ b/lix/libutil/async-io.cc @@ -83,4 +83,54 @@ try { } catch (...) { return {result::current_exception()}; } + +kj::Promise> 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> 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> 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(); +} } diff --git a/lix/libutil/async-io.hh b/lix/libutil/async-io.hh index 6f54c97e2..d71b043b0 100644 --- a/lix/libutil/async-io.hh +++ b/lix/libutil/async-io.hh @@ -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 @@ -81,6 +83,28 @@ public: kj::Promise> read(void * data, size_t len) override; }; +class AsyncBufferedInputStream : public AsyncInputStream +{ + AsyncInputStream & inner; + ref buffer; + +public: + AsyncBufferedInputStream(AsyncInputStream & inner, ref buffer) + : inner(inner) + , buffer(buffer) + { + } + + AsyncBufferedInputStream(AsyncInputStream & inner, size_t bufSize = 32 * 1024) + : AsyncBufferedInputStream(inner, make_ref(bufSize)) + { + } + + KJ_DISALLOW_COPY_AND_MOVE(AsyncBufferedInputStream); + + kj::Promise> 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 buffer; + +public: + AsyncBufferedOutputStream(AsyncOutputStream & inner, ref buffer) + : inner(inner) + , buffer(buffer) + { + } + + AsyncBufferedOutputStream(AsyncOutputStream & inner, size_t bufSize = 32 * 1024) + : AsyncBufferedOutputStream(inner, make_ref(bufSize)) + { + } + + KJ_DISALLOW_COPY_AND_MOVE(AsyncBufferedOutputStream); + + kj::Promise> write(const void * src, size_t size) override; + kj::Promise> flush(); +}; }