diff --git a/doc/manual/rl-next/async-query.md b/doc/manual/rl-next/async-query.md new file mode 100644 index 000000000..b9b14c16c --- /dev/null +++ b/doc/manual/rl-next/async-query.md @@ -0,0 +1,20 @@ +--- +synopsis: "Improved susbtituter query speed" +issues: [] +cls: [] +category: Improvements +credits: [horrors] +--- + +The code used to query substituters for derivations has been rewritten slightly +to take advantage of our asynchronous runtime. Such queries run for every build +that could download from substituters and processes every derivation that isn't +yet present on the local system. Previously Lix would use `http-connections` to +limit query concurrency, even for modern caches that support HTTP/2 and have no +limit on how many queries can be run concurrently on one single connection. Lix +no longer does this, resulting in approximately 60% reduction in query time for +medium-sized closures (e.g. NixOS system closures) during testing, although the +exact number depends greatly on local network latency and generally improves as +latency increases. Unlike previously setting `http-connections` to `1` or other +low values no longer brings a massive penalty in query performance if the cache +in use by the querying system supports HTTP/2 (as e.g. `cache.nixos.org` does). diff --git a/lix/libstore/misc.cc b/lix/libstore/misc.cc index b46383047..9e899a3af 100644 --- a/lix/libstore/misc.cc +++ b/lix/libstore/misc.cc @@ -2,6 +2,7 @@ #include "lix/libstore/parsed-derivations.hh" #include "lix/libstore/globals.hh" #include "lix/libstore/store-api.hh" +#include "lix/libutil/async-collect.hh" #include "lix/libutil/async.hh" #include "lix/libutil/result.hh" #include "lix/libutil/thread-pool.hh" @@ -11,6 +12,7 @@ #include "lix/libutil/strings.hh" #include #include +#include namespace nix { @@ -128,10 +130,7 @@ struct QueryMissingContext DrvState(size_t left) : left(left) { } }; - Sync state_; - - // FIXME: make async. - ThreadPool pool{"queryMissing pool", fileTransferSettings.httpConnections}; + State state; explicit QueryMissingContext( Store & store, @@ -142,7 +141,7 @@ struct QueryMissingContext uint64_t & narSize_ ) : store(store) - , state_{{{}, unknown_, willSubstitute_, willBuild_, downloadSize_, narSize_}} + , state{{}, unknown_, willSubstitute_, willBuild_, downloadSize_, narSize_} { } @@ -150,158 +149,160 @@ struct QueryMissingContext kj::Promise> queryMissing(const std::vector & targets); - void enqueueDerivedPaths(DerivedPathOpaque inputDrv, const StringSet & inputNode) - { + kj::Promise> + enqueueDerivedPaths(DerivedPathOpaque inputDrv, const StringSet & inputNode) + try { if (!inputNode.empty()) { - pool.enqueueWithAio([this, path{DerivedPath::Built{std::move(inputDrv), inputNode}}]( - AsyncIoRoot & aio - ) { doPath(aio, path); }); + TRY_AWAIT(doPath(DerivedPath::Built{std::move(inputDrv), inputNode})); } + co_return result::success(); + } catch (...) { + co_return result::current_exception(); } - void mustBuildDrv(const StorePath & drvPath, const Derivation & drv) - { - { - auto state(state_.lock()); - state->willBuild.insert(drvPath); - } + kj::Promise> mustBuildDrv(const StorePath & drvPath, const Derivation & drv) + try { + state.willBuild.insert(drvPath); - for (const auto & [inputDrv, inputNode] : drv.inputDrvs) { - enqueueDerivedPaths(makeConstantStorePath(inputDrv), inputNode); - } + TRY_AWAIT(asyncSpread(drv.inputDrvs, [&](const auto & input) { + const auto & [inputDrv, inputNode] = input; + return enqueueDerivedPaths(makeConstantStorePath(inputDrv), inputNode); + })); + co_return result::success(); + } catch (...) { + co_return result::current_exception(); } - void checkOutput( - AsyncIoRoot & aio, + kj::Promise> checkOutput( const StorePath & drvPath, ref drv, const StorePath & outPath, - ref> drvState_ + DrvState & drvState ) - { - if (drvState_->lock()->done) return; - + try { SubstitutablePathInfos infos; auto * cap = getDerivationCA(*drv); - aio.blockOn(store.querySubstitutablePathInfos({ + TRY_AWAIT(store.querySubstitutablePathInfos( { - outPath, - cap ? std::optional { *cap } : std::nullopt, + { + outPath, + cap ? std::optional{*cap} : std::nullopt, + }, }, - }, infos)); + infos + )); if (infos.empty()) { - drvState_->lock()->done = true; - mustBuildDrv(drvPath, *drv); + drvState.done = true; + TRY_AWAIT(mustBuildDrv(drvPath, *drv)); } else { { - auto drvState(drvState_->lock()); - if (drvState->done) return; - assert(drvState->left); - drvState->left--; - drvState->outPaths.insert(outPath); - if (!drvState->left) { - for (auto & path : drvState->outPaths) { - pool.enqueueWithAio([this, - path{DerivedPath::Opaque{path}}](AsyncIoRoot & aio) { - doPath(aio, path); - }); - } + if (drvState.done) { + co_return result::success(); + } + assert(drvState.left); + drvState.left--; + drvState.outPaths.insert(outPath); + if (!drvState.left) { + TRY_AWAIT(asyncSpread(drvState.outPaths, [&](auto & path) { + return doPath(DerivedPath::Opaque{path}); + })); } } } + + co_return result::success(); + } catch (...) { + co_return result::current_exception(); } - void doPath(AsyncIoRoot & aio, const DerivedPath & req) + kj::Promise> doPath(DerivedPath req) { - { - auto state(state_.lock()); - if (!state->done.insert(req.to_string(store)).second) return; + if (!state.done.insert(req.to_string(store)).second) { + return {result::success()}; } - std::visit( + return std::visit( overloaded{ - [&](const DerivedPath::Built & bfd) { doPathBuilt(aio, bfd); }, - [&](const DerivedPath::Opaque & bo) { doPathOpaque(aio, bo); }, + [&](DerivedPath::Built bfd) { return doPathBuilt(std::move(bfd)); }, + [&](DerivedPath::Opaque bo) { return doPathOpaque(std::move(bo)); }, }, - req.raw() + std::move(req.raw()) ); } - void doPathBuilt(AsyncIoRoot & aio, const DerivedPath::Built & bfd) - { + kj::Promise> doPathBuilt(DerivedPath::Built bfd) + try { auto & drvPath = bfd.drvPath.path; - if (!aio.blockOn(store.isValidPath(drvPath))) { + if (!TRY_AWAIT(store.isValidPath(drvPath))) { // FIXME: we could try to substitute the derivation. - auto state(state_.lock()); - state->unknown.insert(drvPath); - return; + state.unknown.insert(drvPath); + co_return result::success(); } StorePathSet invalid; - for (auto & [outputName, path] : - aio.blockOn(store.queryDerivationOutputMap(drvPath))) - { - if (bfd.outputs.contains(outputName) && !aio.blockOn(store.isValidPath(path))) + for (auto & [outputName, path] : TRY_AWAIT(store.queryDerivationOutputMap(drvPath))) { + if (bfd.outputs.contains(outputName) && !TRY_AWAIT(store.isValidPath(path))) { invalid.insert(path); + } + } + if (invalid.empty()) { + co_return result::success(); } - if (invalid.empty()) return; - auto drv = make_ref(aio.blockOn(store.derivationFromPath(drvPath))); + auto drv = make_ref(TRY_AWAIT(store.derivationFromPath(drvPath))); ParsedDerivation parsedDrv(StorePath(drvPath), *drv); if (settings.useSubstitutes && parsedDrv.substitutesAllowed()) { - auto drvState = make_ref>(DrvState(invalid.size())); - for (auto & output : invalid) { - pool.enqueueWithAio([=, this](AsyncIoRoot & aio) { - checkOutput(aio, drvPath, drv, output, drvState); - }); - } + DrvState drvState(invalid.size()); + TRY_AWAIT(asyncSpread(invalid, [&](auto & output) { + return checkOutput(drvPath, drv, output, drvState); + })); } else { - mustBuildDrv(drvPath, *drv); + TRY_AWAIT(mustBuildDrv(drvPath, *drv)); } + + co_return result::success(); + } catch (...) { + co_return result::current_exception(); } - void doPathOpaque(AsyncIoRoot & aio, const DerivedPath::Opaque & bo) - { - if (aio.blockOn(store.isValidPath(bo.path))) return; + kj::Promise> doPathOpaque(DerivedPath::Opaque bo) + try { + if (TRY_AWAIT(store.isValidPath(bo.path))) { + co_return result::success(); + } SubstitutablePathInfos infos; - aio.blockOn(store.querySubstitutablePathInfos({{bo.path, std::nullopt}}, infos)); + TRY_AWAIT(store.querySubstitutablePathInfos({{bo.path, std::nullopt}}, infos)); if (infos.empty()) { - auto state(state_.lock()); - state->unknown.insert(bo.path); - return; + state.unknown.insert(bo.path); + co_return result::success(); } auto info = infos.find(bo.path); assert(info != infos.end()); - { - auto state(state_.lock()); - state->willSubstitute.insert(bo.path); - state->downloadSize += info->second.downloadSize; - state->narSize += info->second.narSize; - } + state.willSubstitute.insert(bo.path); + state.downloadSize += info->second.downloadSize; + state.narSize += info->second.narSize; - for (auto & ref : info->second.references) { - pool.enqueueWithAio([this, path{DerivedPath::Opaque{ref}}](AsyncIoRoot & aio) { - doPath(aio, path); - }); - } + TRY_AWAIT(asyncSpread(info->second.references, [&](auto & ref) { + return doPath(DerivedPath::Opaque{ref}); + })); + + co_return result::success(); + } catch (...) { + co_return result::current_exception(); } }; } kj::Promise> QueryMissingContext::queryMissing(const std::vector & targets) try { - for (auto & path : targets) { - pool.enqueueWithAio([=, this](AsyncIoRoot & aio) { doPath(aio, path); }); - } - - TRY_AWAIT(pool.processAsync()); + TRY_AWAIT(asyncSpread(targets, [&](auto & path) { return doPath(path); })); co_return result::success(); } catch (...) { co_return result::current_exception();