diff --git a/lix/libstore/remote-store-connection.hh b/lix/libstore/remote-store-connection.hh index dd255c8b3..2f97cc642 100644 --- a/lix/libstore/remote-store-connection.hh +++ b/lix/libstore/remote-store-connection.hh @@ -7,7 +7,9 @@ #include "lix/libutil/pool.hh" #include "lix/libutil/result.hh" #include "lix/libutil/serialise.hh" +#include #include +#include namespace nix { @@ -140,8 +142,29 @@ struct RemoteStore::ConnectionHandle template kj::Promise> sendCommand(Args &&... args) try { - ((handle->to << std::forward(args)), ...); - processStderr(); + constexpr auto LastArgIdx = sizeof...(Args) - 1; + using AllArgsT = std::tuple; + using LastArgT = std::tuple_element_t; + + // if the last argument can be serialized normally we will serialize *all* + // arguments at once and hand off to the remote. if the last argument does + // not have a serializer we assume it's a callback for a subframe protocol + // and serialize all *preceding* arguments normally before handing over to + // the subframing layer (which is then responsible for any error handling) + if constexpr (requires { handle->to << std::declval(); }) { + ((handle->to << std::forward(args)), ...); + processStderr(); + } else { + using ImmediateArgsIdxs = std::make_index_sequence; + AllArgsT allArgs(std::forward(args)...); + + [&](std::integer_sequence) { + ((handle->to << std::get(std::forward(allArgs))), ...); + }(ImmediateArgsIdxs{}); + + LIX_TRY_AWAIT(withFramedSinkAsync(std::get(allArgs))); + } + if constexpr (std::is_void_v) { co_return result::success(); } else { diff --git a/lix/libstore/remote-store.cc b/lix/libstore/remote-store.cc index 91b181f3c..a4ac8b92a 100644 --- a/lix/libstore/remote-store.cc +++ b/lix/libstore/remote-store.cc @@ -354,24 +354,18 @@ kj::Promise>> RemoteStore::addCAToStore( try { auto conn(TRY_AWAIT(getConnection())); - conn->to - << WorkerProto::Op::AddToStore - << name - << caMethod.render(hashType); - conn->to << WorkerProto::write(*conn, references); - conn->to << repair; - // The dump source may invoke the store, so we need to make some room. connections->incCapacity(); - { - Finally cleanup([&]() { connections->decCapacity(); }); - TRY_AWAIT(conn.withFramedSinkAsync([&](Sink & sink) { - return dump.drainInto(sink); - })); - } + Finally cleanup([&]() { connections->decCapacity(); }); - co_return make_ref( - WorkerProto::Serialise::read(*conn)); + co_return make_ref(TRY_AWAIT(conn.sendCommand( + WorkerProto::Op::AddToStore, + name, + caMethod.render(hashType), + WorkerProto::write(*conn, references), + repair, + [&](Sink & sink) { return dump.drainInto(sink); } + ))); } catch (...) { co_return result::current_exception(); } @@ -401,19 +395,22 @@ kj::Promise> RemoteStore::addToStore( try { auto conn(TRY_AWAIT(getConnection())); - conn->to << WorkerProto::Op::AddToStoreNar - << printStorePath(info.path) - << (info.deriver ? printStorePath(*info.deriver) : "") - << info.narHash.to_string(Base::Base16, false); - conn->to << WorkerProto::write(*conn, info.references); - conn->to << info.registrationTime << info.narSize - << info.ultimate << info.sigs << renderContentAddress(info.ca) - << repair << !checkSigs; - auto copier = copyNAR(source); - TRY_AWAIT(conn.withFramedSinkAsync([&](Sink & sink) { - return copier->drainInto(sink); - })); + TRY_AWAIT(conn.sendCommand( + WorkerProto::Op::AddToStoreNar, + printStorePath(info.path), + (info.deriver ? printStorePath(*info.deriver) : ""), + info.narHash.to_string(Base::Base16, false), + WorkerProto::write(*conn, info.references), + info.registrationTime, + info.narSize, + info.ultimate, + info.sigs, + renderContentAddress(info.ca), + repair, + !checkSigs, + [&](Sink & sink) { return copier->drainInto(sink); } + )); co_return result::success(); } catch (...) { co_return result::current_exception(); @@ -429,25 +426,26 @@ try { auto remoteVersion = TRY_AWAIT(getProtocol()); auto conn(TRY_AWAIT(getConnection())); - conn->to - << WorkerProto::Op::AddMultipleToStore - << repair - << !checkSigs; - // NOLINTNEXTLINE(cppcoreguidelines-avoid-capturing-lambda-coroutines) - TRY_AWAIT(conn.withFramedSinkAsync([&](Sink & sink) -> kj::Promise> { - try { - sink << pathsToCopy.size(); - for (auto & [pathInfo, pathSource] : pathsToCopy) { - sink << WorkerProto::Serialise::write( - WorkerProto::WriteConn {*this, remoteVersion}, - pathInfo); - TRY_AWAIT(TRY_AWAIT(pathSource())->drainInto(sink)); + TRY_AWAIT(conn.sendCommand( + WorkerProto::Op::AddMultipleToStore, + repair, + !checkSigs, + // NOLINTNEXTLINE(cppcoreguidelines-avoid-capturing-lambda-coroutines) + [&](Sink & sink) -> kj::Promise> { + try { + sink << pathsToCopy.size(); + for (auto & [pathInfo, pathSource] : pathsToCopy) { + sink << WorkerProto::Serialise::write( + WorkerProto::WriteConn{*this, remoteVersion}, pathInfo + ); + TRY_AWAIT(TRY_AWAIT(pathSource())->drainInto(sink)); + } + co_return result::success(); + } catch (...) { + co_return result::current_exception(); } - co_return result::success(); - } catch (...) { - co_return result::current_exception(); } - })); + )); co_return result::success(); } catch (...) { co_return result::current_exception(); @@ -654,12 +652,12 @@ try { kj::Promise> RemoteStore::addBuildLog(const StorePath & drvPath, std::string_view log) try { auto conn(TRY_AWAIT(getConnection())); - conn->to << WorkerProto::Op::AddBuildLog << drvPath.to_string(); AsyncStringInputStream source(log); - TRY_AWAIT(conn.withFramedSinkAsync([&](Sink & sink) { - return source.drainInto(sink); - })); - readInt(conn->from); + TRY_AWAIT(conn.sendCommand( + WorkerProto::Op::AddBuildLog, + drvPath.to_string(), + [&](Sink & sink) { return source.drainInto(sink); } + )); co_return result::success(); } catch (...) { co_return result::current_exception();