libstore: asyncify Store::queryRealisation{,Uncached}

Change-Id: I4c4d75c773493fcd63c131e30a2e4d6ec5d230ac
This commit is contained in:
eldritch horrors
2025-03-05 18:49:45 +01:00
parent e057497f3e
commit 03ab0fcb53
19 changed files with 139 additions and 82 deletions
+1 -1
View File
@@ -353,7 +353,7 @@ connected:
for (auto & outputName : wantedOutputs) {
auto thisOutputHash = outputHashes.at(outputName);
auto thisOutputId = DrvOutput{ thisOutputHash, outputName };
if (!store->queryRealisation(thisOutputId)) {
if (!aio.blockOn(store->queryRealisation(thisOutputId))) {
debug("missing output %s", outputName);
assert(optResult);
auto & result = *optResult;
+23 -8
View File
@@ -1,6 +1,8 @@
#include "lix/libcmd/built-path.hh"
#include "lix/libstore/derivations.hh"
#include "lix/libstore/store-api.hh"
#include "lix/libutil/async.hh"
#include "lix/libutil/result.hh"
#include <nlohmann/json.hpp>
@@ -126,10 +128,19 @@ try {
kj::Promise<Result<RealisedPath::Set>> BuiltPath::toRealisedPaths(Store & store) const
try {
RealisedPath::Set res;
std::visit(
overloaded{
[&](const BuiltPath::Opaque & p) { res.insert(p.path); },
[&](const BuiltPath::Built & p) {
auto handlers = overloaded{
// NOLINTNEXTLINE(cppcoreguidelines-avoid-capturing-lambda-coroutines)
[&](const BuiltPath::Opaque & p) -> kj::Promise<Result<void>> {
try {
res.insert(p.path);
return {result::success()};
} catch (...) {
return {result::current_exception()};
}
},
// NOLINTNEXTLINE(cppcoreguidelines-avoid-capturing-lambda-coroutines)
[&](const BuiltPath::Built & p) -> kj::Promise<Result<void>> {
try {
auto drvHashes =
staticOutputHashes(store, store.readDerivation(p.drvPath->outPath()));
for (auto& [outputName, outputPath] : p.outputs) {
@@ -140,8 +151,8 @@ try {
throw Error(
"the derivation '%s' has unrealised output '%s' (derived-path.cc/toRealisedPaths)",
store.printStorePath(p.drvPath->outPath()), outputName);
auto thisRealisation = store.queryRealisation(
DrvOutput{*drvOutput, outputName});
auto thisRealisation = TRY_AWAIT(store.queryRealisation(
DrvOutput{*drvOutput, outputName}));
assert(thisRealisation); // Weve built it, so we must
// have the realisation
res.insert(*thisRealisation);
@@ -149,9 +160,13 @@ try {
res.insert(outputPath);
}
}
},
co_return result::success();
} catch (...) {
co_return result::current_exception();
}
},
raw());
};
TRY_AWAIT(std::visit(handlers, raw()));
co_return res;
} catch (...) {
co_return result::current_exception();
+7 -4
View File
@@ -493,16 +493,19 @@ try {
co_return result::current_exception();
}
std::shared_ptr<const Realisation> BinaryCacheStore::queryRealisationUncached(const DrvOutput & id)
{
kj::Promise<Result<std::shared_ptr<const Realisation>>>
BinaryCacheStore::queryRealisationUncached(const DrvOutput & id)
try {
auto outputInfoFilePath = realisationsPrefix + "/" + id.to_string() + ".doi";
auto data = getFileContents(outputInfoFilePath);
if (!data) return {};
if (!data) co_return result::success(nullptr);
auto realisation = Realisation::fromJSON(
nlohmann::json::parse(*data), outputInfoFilePath);
return std::make_shared<const Realisation>(realisation);
co_return std::make_shared<const Realisation>(realisation);
} catch (...) {
co_return result::current_exception();
}
kj::Promise<Result<void>> BinaryCacheStore::registerDrvOutput(const Realisation& info)
+2 -1
View File
@@ -147,7 +147,8 @@ public:
kj::Promise<Result<void>> registerDrvOutput(const Realisation & info) override;
std::shared_ptr<const Realisation> queryRealisationUncached(const DrvOutput &) override;
kj::Promise<Result<std::shared_ptr<const Realisation>>>
queryRealisationUncached(const DrvOutput &) override;
box_ptr<Source> narFromPath(const StorePath & path) override;
+21 -14
View File
@@ -1163,21 +1163,28 @@ try {
"derivation '%s' doesn't have expected output '%s' (derivation-goal.cc/resolvedFinished,resolve)",
worker.store.printStorePath(drvPath), outputName);
auto realisation = [&]{
auto take1 = get(resolvedResult.builtOutputs, outputName);
if (take1) return *take1;
// NOLINTNEXTLINE(cppcoreguidelines-avoid-capturing-lambda-coroutines)
auto realisation = TRY_AWAIT([&]() -> kj::Promise<Result<Realisation>> {
try {
auto take1 = get(resolvedResult.builtOutputs, outputName);
if (take1) co_return *take1;
/* The above `get` should work. But sateful tracking of
outputs in resolvedResult, this can get out of sync with the
store, which is our actual source of truth. For now we just
check the store directly if it fails. */
auto take2 = worker.evalStore.queryRealisation(DrvOutput { *resolvedHash, outputName });
if (take2) return *take2;
/* The above `get` should work. But sateful tracking of
outputs in resolvedResult, this can get out of sync with the
store, which is our actual source of truth. For now we just
check the store directly if it fails. */
auto take2 = TRY_AWAIT(
worker.evalStore.queryRealisation(DrvOutput{*resolvedHash, outputName})
);
if (take2) co_return *take2;
throw Error(
"derivation '%s' doesn't have expected output '%s' (derivation-goal.cc/resolvedFinished,realisation)",
worker.store.printStorePath(resolvedDrvGoal->drvPath), outputName);
}();
throw Error(
"derivation '%s' doesn't have expected output '%s' (derivation-goal.cc/resolvedFinished,realisation)",
worker.store.printStorePath(resolvedDrvGoal->drvPath), outputName);
} catch (...) {
co_return result::current_exception();
}
}());
if (drv->type().isPure()) {
auto newRealisation = realisation;
@@ -1681,7 +1688,7 @@ try {
}
auto drvOutput = DrvOutput{info.outputHash, i.first};
if (experimentalFeatureSettings.isEnabled(Xp::CaDerivations)) {
if (auto real = worker.store.queryRealisation(drvOutput)) {
if (auto real = TRY_AWAIT(worker.store.queryRealisation(drvOutput))) {
info.known = {
.path = real->outPath,
.status = PathStatus::Valid,
@@ -30,7 +30,7 @@ try {
trace("init");
/* If the derivation already exists, were done */
if (worker.store.queryRealisation(id)) {
if (TRY_AWAIT(worker.store.queryRealisation(id))) {
co_return WorkResult{ecSuccess};
}
@@ -79,7 +79,8 @@ try {
std::async(std::launch::async, [downloadState{downloadState}, id{id}, sub{sub}] {
Finally updateStats([&]() { downloadState->outPipe->fulfill(); });
ReceiveInterrupts receiveInterrupts;
return sub->queryRealisation(id);
AsyncIoRoot aio;
return aio.blockOn(sub->queryRealisation(id));
});
co_await pipe.promise;
@@ -107,7 +108,7 @@ try {
kj::Vector<std::pair<GoalPtr, kj::Promise<Result<WorkResult>>>> dependencies;
for (const auto & [depId, depPath] : outputInfo->dependentRealisations) {
if (depId != id) {
if (auto localOutputInfo = worker.store.queryRealisation(depId);
if (auto localOutputInfo = TRY_AWAIT(worker.store.queryRealisation(depId));
localOutputInfo && localOutputInfo->outPath != depPath) {
warn(
"substituter '%s' has an incompatible realisation for '%s', ignoring.\n"
+7 -4
View File
@@ -1164,13 +1164,16 @@ struct RestrictedStore : public virtual IndirectRootStore, public virtual GcStor
// corresponds to an allowed derivation
try { throw Error("registerDrvOutput"); } catch (...) { return {result::current_exception()}; }
std::shared_ptr<const Realisation> queryRealisationUncached(const DrvOutput & id) override
kj::Promise<Result<std::shared_ptr<const Realisation>>>
queryRealisationUncached(const DrvOutput & id) override
// XXX: This should probably be allowed if the realisation corresponds to
// an allowed derivation
{
try {
if (!goal.isAllowed(id))
return nullptr;
return next->queryRealisation(id);
co_return result::success(nullptr);
co_return TRY_AWAIT(next->queryRealisation(id));
} catch (...) {
co_return result::current_exception();
}
kj::Promise<Result<void>> buildPaths(
+1 -1
View File
@@ -969,7 +969,7 @@ static void performOp(AsyncIoRoot & aio, TunnelLogger * logger, ref<Store> store
case WorkerProto::Op::QueryRealisation: {
logger->startWork();
auto outputId = DrvOutput::parse(readString(from));
auto info = store->queryRealisation(outputId);
auto info = aio.blockOn(store->queryRealisation(outputId));
logger->stopWork();
if (GET_PROTOCOL_MINOR(clientVersion) < 31) {
std::set<StorePath> outPaths;
+3 -2
View File
@@ -73,8 +73,9 @@ struct DummyStore final : public Store
box_ptr<Source> narFromPath(const StorePath & path) override
{ unsupported("narFromPath"); }
std::shared_ptr<const Realisation> queryRealisationUncached(const DrvOutput &) override
{ return nullptr; }
kj::Promise<Result<std::shared_ptr<const Realisation>>>
queryRealisationUncached(const DrvOutput &) override
{ co_return result::success(nullptr); }
virtual ref<FSAccessor> getFSAccessor() override
{ unsupported("getFSAccessor"); }
+3 -2
View File
@@ -461,9 +461,10 @@ public:
return {result::success(std::nullopt)};
}
std::shared_ptr<const Realisation> queryRealisationUncached(const DrvOutput &) override
kj::Promise<Result<std::shared_ptr<const Realisation>>>
queryRealisationUncached(const DrvOutput &) override
// TODO: Implement
{ unsupported("queryRealisation"); }
try { unsupported("queryRealisation"); } catch (...) { co_return result::current_exception(); }
};
void registerLegacySSHStore() {
+17 -8
View File
@@ -1950,16 +1950,25 @@ std::optional<const Realisation> LocalStore::queryRealisation_(
return { res };
}
std::shared_ptr<const Realisation> LocalStore::queryRealisationUncached(const DrvOutput & id)
{
auto maybeRealisation = retrySQLite([&]() {
auto state = dbPool.get();
return queryRealisation_(*state, id);
}, always_progresses);
kj::Promise<Result<std::shared_ptr<const Realisation>>>
LocalStore::queryRealisationUncached(const DrvOutput & id)
try {
auto maybeRealisation =
// NOLINTNEXTLINE(cppcoreguidelines-avoid-capturing-lambda-coroutines)
TRY_AWAIT(retrySQLite([&]() -> kj::Promise<Result<std::optional<const Realisation>>> {
try {
auto state = dbPool.get();
co_return queryRealisation_(*state, id);
} catch (...) {
co_return result::current_exception();
}
}));
if (maybeRealisation)
return std::make_shared<const Realisation>(maybeRealisation.value());
co_return std::make_shared<const Realisation>(maybeRealisation.value());
else
return nullptr;
co_return result::success(nullptr);
} catch (...) {
co_return result::current_exception();
}
ContentAddress LocalStore::hashCAPath(
+2 -1
View File
@@ -332,7 +332,8 @@ public:
std::optional<const Realisation> queryRealisation_(DBState & state, const DrvOutput & id);
std::optional<std::pair<int64_t, Realisation>> queryRealisationCore_(DBState & state, const DrvOutput & id);
std::shared_ptr<const Realisation> queryRealisationUncached(const DrvOutput&) override;
kj::Promise<Result<std::shared_ptr<const Realisation>>>
queryRealisationUncached(const DrvOutput&) override;
kj::Promise<Result<std::optional<std::string>>> getVersion() override;
+3 -3
View File
@@ -276,7 +276,7 @@ struct QueryMissingContext
bool found = false;
for (auto &sub : aio.blockOn(getDefaultSubstituters())) {
auto realisation = sub->queryRealisation({hash, outputName});
auto realisation = aio.blockOn(sub->queryRealisation({hash, outputName}));
if (!realisation)
continue;
found = true;
@@ -432,8 +432,8 @@ try {
throw Error(
"output '%s' of derivation '%s' isn't realised", outputName,
store.printStorePath(inputDrv));
auto thisRealisation = store.queryRealisation(
DrvOutput{*outputHash, outputName});
auto thisRealisation = TRY_AWAIT(store.queryRealisation(
DrvOutput{*outputHash, outputName}));
if (!thisRealisation)
throw Error(
"output '%s' of derivation '%s' isnt built", outputName,
+16 -10
View File
@@ -1,5 +1,6 @@
#include "lix/libstore/realisation.hh"
#include "lix/libstore/store-api.hh"
#include "lix/libutil/async.hh"
#include "lix/libutil/closure.hh"
#include "lix/libutil/result.hh"
#include <nlohmann/json.hpp>
@@ -37,19 +38,24 @@ kj::Promise<Result<void>> Realisation::closure(
Store & store, const std::set<Realisation> & startOutputs, std::set<Realisation> & res
)
try {
auto getDeps = [&](const Realisation& current) -> std::set<Realisation> {
std::set<Realisation> res;
for (auto& [currentDep, _] : current.dependentRealisations) {
if (auto currentRealisation = store.queryRealisation(currentDep))
res.insert(*currentRealisation);
else
throw Error(
"Unrealised derivation '%s'", currentDep.to_string());
// NOLINTNEXTLINE(cppcoreguidelines-avoid-capturing-lambda-coroutines)
auto getDeps = [&](const Realisation& current) -> kj::Promise<Result<std::set<Realisation>>> {
try {
std::set<Realisation> res;
for (auto& [currentDep, _] : current.dependentRealisations) {
if (auto currentRealisation = TRY_AWAIT(store.queryRealisation(currentDep)))
res.insert(*currentRealisation);
else
throw Error(
"Unrealised derivation '%s'", currentDep.to_string());
}
co_return res;
} catch (...) {
co_return result::current_exception();
}
return res;
};
res.merge(computeClosure<Realisation>(startOutputs, getDeps));
res.merge(TRY_AWAIT(computeClosureAsync<Realisation>(startOutputs, getDeps)));
co_return result::success();
} catch (...) {
co_return result::current_exception();
+13 -8
View File
@@ -628,13 +628,14 @@ try {
co_return result::current_exception();
}
std::shared_ptr<const Realisation> RemoteStore::queryRealisationUncached(const DrvOutput & id)
{
kj::Promise<Result<std::shared_ptr<const Realisation>>>
RemoteStore::queryRealisationUncached(const DrvOutput & id)
try {
auto conn(getConnection());
if (GET_PROTOCOL_MINOR(conn->daemonVersion) < 27) {
warn("the daemon is too old to support content-addressed derivations, please upgrade it to 2.4");
return nullptr;
co_return result::success(nullptr);
}
conn->to << WorkerProto::Op::QueryRealisation;
@@ -645,15 +646,19 @@ std::shared_ptr<const Realisation> RemoteStore::queryRealisationUncached(const D
auto outPaths = WorkerProto::Serialise<std::set<StorePath>>::read(
*this, *conn);
if (outPaths.empty())
return nullptr;
return std::make_shared<const Realisation>(Realisation { .id = id, .outPath = *outPaths.begin() });
co_return result::success(nullptr);
co_return std::make_shared<const Realisation>(
Realisation{.id = id, .outPath = *outPaths.begin()}
);
} else {
auto realisations = WorkerProto::Serialise<std::set<Realisation>>::read(
*this, *conn);
if (realisations.empty())
return nullptr;
return std::make_shared<const Realisation>(*realisations.begin());
co_return result::success(nullptr);
co_return std::make_shared<const Realisation>(*realisations.begin());
}
} catch (...) {
co_return result::current_exception();
}
kj::Promise<Result<void>> RemoteStore::copyDrvsFromEvalStore(
@@ -764,7 +769,7 @@ try {
auto outputId = DrvOutput{ *outputHash, output };
if (experimentalFeatureSettings.isEnabled(Xp::CaDerivations)) {
auto realisation =
queryRealisation(outputId);
TRY_AWAIT(queryRealisation(outputId));
if (!realisation)
throw MissingRealisation(outputId);
res.builtOutputs.emplace(output, *realisation);
+2 -1
View File
@@ -116,7 +116,8 @@ public:
kj::Promise<Result<void>> registerDrvOutput(const Realisation & info) override;
std::shared_ptr<const Realisation> queryRealisationUncached(const DrvOutput &) override;
kj::Promise<Result<std::shared_ptr<const Realisation>>>
queryRealisationUncached(const DrvOutput &) override;
kj ::Promise<Result<void>> buildPaths(
const std::vector<DerivedPath> & paths,
+10 -8
View File
@@ -532,7 +532,7 @@ try {
auto drv = evalStore.readInvalidDerivation(path);
auto drvHashes = staticOutputHashes(*this, drv);
for (auto & [outputName, hash] : drvHashes) {
auto realisation = queryRealisation(DrvOutput{hash, outputName});
auto realisation = TRY_AWAIT(queryRealisation(DrvOutput{hash, outputName}));
if (realisation) {
outputs.insert_or_assign(outputName, realisation->outPath);
} else {
@@ -746,8 +746,8 @@ ref<const ValidPathInfo> Store::queryPathInfo(const StorePath & storePath)
return ref<const ValidPathInfo>(info);
}
std::shared_ptr<const Realisation> Store::queryRealisation(const DrvOutput & id)
{
kj::Promise<Result<std::shared_ptr<const Realisation>>> Store::queryRealisation(const DrvOutput & id)
try {
if (diskCache) {
auto [cacheOutcome, maybeCachedRealisation]
@@ -755,18 +755,18 @@ std::shared_ptr<const Realisation> Store::queryRealisation(const DrvOutput & id)
switch (cacheOutcome) {
case NarInfoDiskCache::oValid:
debug("Returning a cached realisation for %s", id.to_string());
return maybeCachedRealisation;
co_return maybeCachedRealisation;
case NarInfoDiskCache::oInvalid:
debug(
"Returning a cached missing realisation for %s",
id.to_string());
return nullptr;
co_return result::success(nullptr);
case NarInfoDiskCache::oUnknown:
break;
}
}
auto info = queryRealisationUncached(id);
auto info = TRY_AWAIT(queryRealisationUncached(id));
if (diskCache) {
if (info)
@@ -775,7 +775,9 @@ std::shared_ptr<const Realisation> Store::queryRealisation(const DrvOutput & id)
diskCache->upsertAbsentRealisation(getUri(), id);
}
return info;
co_return info;
} catch (...) {
co_return result::current_exception();
}
kj::Promise<Result<void>> Store::substitutePaths(const StorePathSet & paths)
@@ -1161,7 +1163,7 @@ try {
[&](AsyncIoRoot & aio, const Realisation & current) -> std::set<Realisation> {
std::set<Realisation> children;
for (const auto & [drvOutput, _] : current.dependentRealisations) {
auto currentChild = srcStore.queryRealisation(drvOutput);
auto currentChild = aio.blockOn(srcStore.queryRealisation(drvOutput));
if (!currentChild)
throw Error(
"incomplete realisation closure: '%s' is a "
+3 -2
View File
@@ -390,7 +390,7 @@ public:
/**
* Query the information about a realisation.
*/
std::shared_ptr<const Realisation> queryRealisation(const DrvOutput &);
kj::Promise<Result<std::shared_ptr<const Realisation>>> queryRealisation(const DrvOutput &);
/**
@@ -421,7 +421,8 @@ protected:
* Note to implementors: should return `nullptr` when the path is not found.
*/
virtual std::shared_ptr<const ValidPathInfo> queryPathInfoUncached(const StorePath & path) = 0;
virtual std::shared_ptr<const Realisation> queryRealisationUncached(const DrvOutput &) = 0;
virtual kj::Promise<Result<std::shared_ptr<const Realisation>>>
queryRealisationUncached(const DrvOutput &) = 0;
public:
+1 -1
View File
@@ -134,7 +134,7 @@ SV * queryPathInfo(char * path, int base32)
SV * queryRawRealisation(char * outputId)
PPCODE:
try {
auto realisation = store()->queryRealisation(DrvOutput::parse(outputId));
auto realisation = aio().blockOn(store()->queryRealisation(DrvOutput::parse(outputId)));
if (realisation)
XPUSHs(sv_2mortal(newSVpv(realisation->toJSON().dump().c_str(), 0)));
else