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
This commit is contained in:
+13
-13
@@ -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;
|
||||
}
|
||||
|
||||
|
||||
|
||||
@@ -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<IoBuffer> buffer;
|
||||
|
||||
BufferedSink(size_t bufSize = 32 * 1024) : buffer(bufSize) {}
|
||||
BufferedSink(size_t bufSize = 32 * 1024) : buffer(make_ref<IoBuffer>(bufSize)) {}
|
||||
explicit BufferedSink(ref<IoBuffer> buffer) : buffer(std::move(buffer)) {}
|
||||
|
||||
void operator () (std::string_view data) override;
|
||||
|
||||
@@ -97,9 +99,10 @@ struct Source
|
||||
*/
|
||||
struct BufferedSource : Source
|
||||
{
|
||||
IoBuffer buffer;
|
||||
ref<IoBuffer> buffer;
|
||||
|
||||
BufferedSource(size_t bufSize = 32 * 1024) : buffer(bufSize) {}
|
||||
BufferedSource(size_t bufSize = 32 * 1024) : buffer(make_ref<IoBuffer>(bufSize)) {}
|
||||
explicit BufferedSource(ref<IoBuffer> buffer) : buffer(std::move(buffer)) {}
|
||||
|
||||
size_t read(char * data, size_t len) override;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user