From b6e102156a42465e2d98c25609c18aa2edaf1cf2 Mon Sep 17 00:00:00 2001 From: Eelco Dolstra Date: Tue, 7 Jul 2026 15:25:05 +0200 Subject: [PATCH 1/3] Add Store::queryPathInfos() method This is intended to allow a store to query multiple paths in a single call. Taken from https://github.com/DeterminateSystems/nix-src/pull/523. Assisted-by: Claude Fable 5 --- src/libstore/include/nix/store/store-api.hh | 15 +++++++++++++++ src/libstore/store-api.cc | 19 +++++++++++++++++++ 2 files changed, 34 insertions(+) diff --git a/src/libstore/include/nix/store/store-api.hh b/src/libstore/include/nix/store/store-api.hh index 87a0d0f37530..6df804515448 100644 --- a/src/libstore/include/nix/store/store-api.hh +++ b/src/libstore/include/nix/store/store-api.hh @@ -14,6 +14,8 @@ #include "nix/store/store-dir-config.hh" #include "nix/store/store-reference.hh" #include "nix/util/source-path.hh" +#include "nix/util/async.hh" +#include "nix/util/fun.hh" #include #include @@ -430,6 +432,19 @@ public: */ void queryPathInfo(const StorePath & path, Callback> callback) noexcept; + /** + * Asynchronously query information about multiple store paths. As + * results arrive (possibly in batches from a remote server), + * `callback` is invoked one or more times with a vector of + * `(path, info)` pairs. A null `info` denotes that the path is + * not valid. Every requested path is reported exactly once across + * all invocations of `callback`. Unlike `queryPathInfo()`, an + * invalid path is not an error. + */ + virtual asio::awaitable queryPathInfos( + const std::set & paths, + fun>>)> callback); + /** * Version of queryPathInfo() that only queries the local narinfo cache and not * the actual store. diff --git a/src/libstore/store-api.cc b/src/libstore/store-api.cc index d2effb6353cd..d926fa00f23c 100644 --- a/src/libstore/store-api.cc +++ b/src/libstore/store-api.cc @@ -613,6 +613,25 @@ void Store::queryPathInfo(const StorePath & storePath, Callback Store::queryPathInfos( + const std::set & paths, + fun>>)> callback) +{ + /* Default implementation: query each path individually, reporting + each result as it arrives. */ + co_await forEachAsync(paths, [&](const StorePath & path) -> asio::awaitable { + std::shared_ptr info; + try { + auto i = co_await callbackToAwaitable>( + [&](Callback> cb) { queryPathInfo(path, std::move(cb)); }); + info = i.get_ptr(); + } catch (InvalidPath &) { + } + std::vector>> result{{path, info}}; + callback(std::move(result)); + }); +} + void Store::queryRealisation( const DrvOutput & id, Callback> callback) noexcept { From f65f632605ee59670079985fcfd4005e62573753 Mon Sep 17 00:00:00 2001 From: Eelco Dolstra Date: Tue, 7 Jul 2026 16:16:55 +0200 Subject: [PATCH 2/3] libstore: Add a QueryPathInfos daemon protocol operation Add a batched QueryPathInfos operation to the worker protocol, gated behind a new "queryPathInfos" protocol feature, and use it to implement an efficient RemoteStore::queryPathInfos() override. Previously, querying path infos from a remote store required a network round-trip per store path, which is extremely slow over high-latency links (e.g. in computeFSClosure()). Now the entire set of paths is sent in a single operation. The response contains the infos for the valid paths; paths not reported are invalid. The client answers what it can from the in-memory path info cache first, updates the cache with the results (including negative entries), and falls back to the per-path base implementation for daemons that don't support the new operation. Assisted-by: Claude Fable 5 --- src/libstore/daemon.cc | 32 +++++++++ .../include/nix/store/remote-store.hh | 4 ++ .../include/nix/store/worker-protocol.hh | 2 + src/libstore/remote-store.cc | 65 +++++++++++++++++++ src/libstore/worker-protocol.cc | 1 + 5 files changed, 104 insertions(+) diff --git a/src/libstore/daemon.cc b/src/libstore/daemon.cc index 31cc2c512719..e3f18db0cdfc 100644 --- a/src/libstore/daemon.cc +++ b/src/libstore/daemon.cc @@ -21,6 +21,7 @@ #include "nix/store/globals.hh" #include "nix/store/active-builds.hh" #include "nix/util/provenance.hh" +#include "nix/util/async.hh" #ifndef _WIN32 // TODO need graceful async exit support on Windows? # include "nix/util/monitor-fd.hh" @@ -886,6 +887,37 @@ static void performOp( break; } + case WorkerProto::Op::QueryPathInfos: { + auto paths = WorkerProto::Serialise::read(*store, rconn); + logger->startWork(); + std::vector infos; + { + asio::io_context ctx; + std::exception_ptr ex; + asio::co_spawn( + ctx, + [&]() -> asio::awaitable { + co_await store->queryPathInfos( + paths, [&](std::vector>> results) { + for (auto & [path, info] : results) + if (info) + infos.push_back(*info); + }); + }, + [&](std::exception_ptr e) { ex = e; }); + ctx.run(); + if (ex) + std::rethrow_exception(ex); + } + logger->stopWork(); + /* Write the infos for the valid paths. Paths not reported are + invalid. */ + conn.to << infos.size(); + for (auto & info : infos) + WorkerProto::write(*store, wconn, info); + break; + } + case WorkerProto::Op::OptimiseStore: logger->startWork(); store->optimiseStore(); diff --git a/src/libstore/include/nix/store/remote-store.hh b/src/libstore/include/nix/store/remote-store.hh index 3333c7130779..7f3ae5d75b41 100644 --- a/src/libstore/include/nix/store/remote-store.hh +++ b/src/libstore/include/nix/store/remote-store.hh @@ -62,6 +62,10 @@ struct RemoteStore : public virtual Store, void queryPathInfoUncached( const StorePath & path, Callback> callback) noexcept override; + asio::awaitable queryPathInfos( + const std::set & paths, + fun>>)> callback) override; + void queryReferrers(const StorePath & path, StorePathSet & referrers) override; StorePathSet queryValidDerivers(const StorePath & path) override; diff --git a/src/libstore/include/nix/store/worker-protocol.hh b/src/libstore/include/nix/store/worker-protocol.hh index 05b6ecc76499..44c680660fa4 100644 --- a/src/libstore/include/nix/store/worker-protocol.hh +++ b/src/libstore/include/nix/store/worker-protocol.hh @@ -114,6 +114,7 @@ struct WorkerProto static constexpr std::string_view featureProvenance = "provenance"; static constexpr std::string_view featureVersionedAddToStoreMultiple = "versionedAddToStoreMultiple"; static constexpr std::string_view featureAddTempRoots = "addTempRoots"; + static constexpr std::string_view featureQueryPathInfos = "queryPathInfos"; /** * A unidirectional read connection, to be used by the read half of the @@ -240,6 +241,7 @@ enum struct WorkerProto::Op : uint64_t { AddPermRoot = 47, QueryActiveBuilds = 48, AddTempRoots = 49, + QueryPathInfos = 50, }; struct WorkerProto::ClientHandshakeInfo diff --git a/src/libstore/remote-store.cc b/src/libstore/remote-store.cc index 0532029f872c..47cdba6551f2 100644 --- a/src/libstore/remote-store.cc +++ b/src/libstore/remote-store.cc @@ -246,6 +246,71 @@ void RemoteStore::queryPathInfoUncached( } } +asio::awaitable RemoteStore::queryPathInfos( + const std::set & paths, + fun>>)> callback) +{ + /* Filter out paths that we already have cached. */ + StorePathSet uncached; + { + std::vector>> cached; + for (auto & path : paths) { + if (auto r = queryPathInfoFromClientCache(path)) + cached.emplace_back(path, *r); + else + uncached.insert(path); + } + if (!cached.empty()) + callback(std::move(cached)); + } + + if (uncached.empty()) + co_return; + + { + auto conn(getConnection()); + + if (conn->protoVersion.features.contains(WorkerProto::featureQueryPathInfos)) { + auto cacheResult = [&](const StorePath & path, std::shared_ptr info) { + pathInfoCache->lock()->upsert(path, PathInfoCacheValue{.value = info}); + }; + + conn->to << WorkerProto::Op::QueryPathInfos; + WorkerProto::write(*this, *conn, uncached); + conn.processStderr(); + + /* Read the infos for the valid paths. Paths not reported + by the daemon are invalid. */ + auto todo = uncached; + auto n = readNum(conn->from); + std::vector>> results; + for (size_t i = 0; i < n; i++) { + auto info = + std::make_shared(WorkerProto::Serialise::read(*this, *conn)); + if (!todo.erase(info->path)) + throw Error("daemon returned path info for unexpected path '%s'", printStorePath(info->path)); + cacheResult(info->path, info); + results.emplace_back(info->path, info); + } + + for (auto & path : todo) { + cacheResult(path, nullptr); + results.emplace_back(path, nullptr); + } + + callback(std::move(results)); + + co_return; + } + + /* Release the connection to prevent the fallback below from + deadlocking on the connection pool. */ + } + + /* Fallback for daemons that don't support the batched operation. */ + co_await Store::queryPathInfos(uncached, std::move(callback)); +} + void RemoteStore::queryReferrers(const StorePath & path, StorePathSet & referrers) { auto conn(getConnection()); diff --git a/src/libstore/worker-protocol.cc b/src/libstore/worker-protocol.cc index 8baab6442317..54f3a8a1c7f0 100644 --- a/src/libstore/worker-protocol.cc +++ b/src/libstore/worker-protocol.cc @@ -27,6 +27,7 @@ const WorkerProto::Version WorkerProto::latest = { std::string{WorkerProto::featureProvenance}, std::string{WorkerProto::featureVersionedAddToStoreMultiple}, std::string{WorkerProto::featureAddTempRoots}, + std::string{WorkerProto::featureQueryPathInfos}, }, }; From ac149f57794c382eec219044ebe2f6e6548666dd Mon Sep 17 00:00:00 2001 From: Eelco Dolstra Date: Tue, 7 Jul 2026 15:53:41 +0200 Subject: [PATCH 3/3] computeFSClosure(): Use queryPathInfos() --- src/libstore/misc.cc | 81 ++++++++++++++++++++++++++++++-------------- 1 file changed, 56 insertions(+), 25 deletions(-) diff --git a/src/libstore/misc.cc b/src/libstore/misc.cc index c6d1a10a95a6..50464bb9b240 100644 --- a/src/libstore/misc.cc +++ b/src/libstore/misc.cc @@ -19,15 +19,11 @@ namespace nix { void Store::computeFSClosure( - const StorePathSet & startPaths, - StorePathSet & paths_, - bool flipDirection, - bool includeOutputs, - bool includeDerivers) + const StorePathSet & startPaths, StorePathSet & out, bool flipDirection, bool includeOutputs, bool includeDerivers) { - std::function(const StorePath & path)> queryDeps; - if (flipDirection) - queryDeps = [this, includeOutputs, includeDerivers](const StorePath & path) -> asio::awaitable { + if (flipDirection) { + std::function(const StorePath & path)> queryDeps = + [this, includeOutputs, includeDerivers](const StorePath & path) -> asio::awaitable { StorePathSet res; StorePathSet referrers; queryReferrers(path, referrers); @@ -45,27 +41,62 @@ void Store::computeFSClosure( res.insert(*maybeOutPath); co_return res; }; - else - queryDeps = [this, includeOutputs, includeDerivers](const StorePath & path) -> asio::awaitable { - StorePathSet res; - auto info = co_await callbackToAwaitable>( - [this, path](Callback> cb) { queryPathInfo(path, std::move(cb)); }); + computeClosure(startPaths, out, GetEdgesAsync(queryDeps)); + } else { + + asio::io_context ctx; + std::exception_ptr ex; + + asio::co_spawn( + ctx, + [&]() -> asio::awaitable { + auto todo = startPaths; + + StorePathSet required, done; + + while (!todo.empty()) { + StorePathSet batch; + for (auto & path : std::exchange(todo, {})) + if (done.insert(path).second) + batch.insert(path); + + co_await queryPathInfos( + batch, [&](std::vector>> infos) { + for (auto & [path, info] : infos) { + if (!info) { + if (required.contains(path)) + throw InvalidPath("path '%s' is not valid", printStorePath(path)); + continue; + } - for (auto & ref : info->references) - if (ref != path) - res.insert(ref); + out.insert(path); - if (includeOutputs && path.isDerivation()) - for (auto & [_, maybeOutPath] : queryPartialDerivationOutputMap(path)) - if (maybeOutPath && isValidPath(*maybeOutPath)) - res.insert(*maybeOutPath); + for (auto & ref : info->references) + if (ref != path) { + required.insert(ref); + todo.insert(ref); + } - if (includeDerivers && info->deriver && isValidPath(*info->deriver)) - res.insert(*info->deriver); - co_return res; - }; + if (includeOutputs && path.isDerivation()) + // FIXME: need an async, multiple-path version of queryPartialDerivationOutputMap(). + for (auto & [_, maybeOutPath] : queryPartialDerivationOutputMap(path)) + if (maybeOutPath) + todo.insert(*maybeOutPath); + + if (includeDerivers && info->deriver) + todo.insert(*info->deriver); + } + }); + } + + co_return; + }, + [&](std::exception_ptr e) { ex = e; }); - computeClosure(startPaths, paths_, GetEdgesAsync(queryDeps)); + ctx.run(); + if (ex) + std::rethrow_exception(ex); + } } void Store::computeFSClosure(