From 56847dc10d52066dde03da01eee1659cdb691556 Mon Sep 17 00:00:00 2001 From: eldritch horrors Date: Wed, 11 Jun 2025 17:22:33 +0200 Subject: [PATCH] libutil: make Buffered{Sink,Source} io buffer shareable async io for remote store connections needs some sync parts still for serialization purposes, and those will have to reuse async io buffers Change-Id: I05e066e3bf8c4318dc23306383f6a849d018ef91 --- lix/libutil/serialise.cc | 26 +++++++++++++------------- lix/libutil/serialise.hh | 11 +++++++---- 2 files changed, 20 insertions(+), 17 deletions(-) diff --git a/lix/libutil/serialise.cc b/lix/libutil/serialise.cc index a97087f4b..125be13c4 100644 --- a/lix/libutil/serialise.cc +++ b/lix/libutil/serialise.cc @@ -53,19 +53,19 @@ void BufferedSink::operator () (std::string_view data) while (!data.empty()) { /* Optimisation: bypass the buffer if the data exceeds the buffer size. */ - if (buffer.used() + data.size() >= buffer.size()) { + if (buffer->used() + data.size() >= buffer->size()) { flush(); writeUnbuffered(data); break; } /* Otherwise, copy the bytes to the buffer. Flush the buffer when it's full. */ - auto into = buffer.getWriteBuffer(); + auto into = buffer->getWriteBuffer(); size_t n = std::min(data.size(), into.size()); memcpy(into.data(), data.data(), n); data.remove_prefix(n); - buffer.added(n); - if (buffer.used() == buffer.size()) { + buffer->added(n); + if (buffer->used() == buffer->size()) { flush(); } } @@ -73,10 +73,10 @@ void BufferedSink::operator () (std::string_view data) void BufferedSink::flush() { - if (buffer.used() > 0) { - auto from = buffer.getReadBuffer(); + if (buffer->used() > 0) { + auto from = buffer->getReadBuffer(); writeUnbuffered({from.data(), from.size()}); - buffer.consumed(from.size()); + buffer->consumed(from.size()); } } @@ -138,22 +138,22 @@ std::string Source::drain() size_t BufferedSource::read(char * data, size_t len) { - if (buffer.used() == 0) { - auto into = buffer.getWriteBuffer(); - buffer.added(readUnbuffered(into.data(), into.size())); + if (buffer->used() == 0) { + auto into = buffer->getWriteBuffer(); + buffer->added(readUnbuffered(into.data(), into.size())); } - auto from = buffer.getReadBuffer(); + auto from = buffer->getReadBuffer(); len = std::min(len, from.size()); memcpy(data, from.data(), len); - buffer.consumed(len); + buffer->consumed(len); return len; } bool BufferedSource::hasData() { - return buffer.used() > 0; + return buffer->used() > 0; } diff --git a/lix/libutil/serialise.hh b/lix/libutil/serialise.hh index 574510d75..c5e46d895 100644 --- a/lix/libutil/serialise.hh +++ b/lix/libutil/serialise.hh @@ -6,6 +6,7 @@ #include "lix/libutil/charptr-cast.hh" #include "lix/libutil/generator.hh" #include "lix/libutil/io-buffer.hh" +#include "lix/libutil/ref.hh" #include "lix/libutil/types.hh" #include "lix/libutil/file-descriptor.hh" @@ -44,9 +45,10 @@ struct FinishSink : virtual Sink */ struct BufferedSink : virtual Sink { - IoBuffer buffer; + ref buffer; - BufferedSink(size_t bufSize = 32 * 1024) : buffer(bufSize) {} + BufferedSink(size_t bufSize = 32 * 1024) : buffer(make_ref(bufSize)) {} + explicit BufferedSink(ref buffer) : buffer(std::move(buffer)) {} void operator () (std::string_view data) override; @@ -97,9 +99,10 @@ struct Source */ struct BufferedSource : Source { - IoBuffer buffer; + ref buffer; - BufferedSource(size_t bufSize = 32 * 1024) : buffer(bufSize) {} + BufferedSource(size_t bufSize = 32 * 1024) : buffer(make_ref(bufSize)) {} + explicit BufferedSource(ref buffer) : buffer(std::move(buffer)) {} size_t read(char * data, size_t len) override;