Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 15 additions & 1 deletion src/libstore/daemon.cc
Original file line number Diff line number Diff line change
Expand Up @@ -515,7 +515,21 @@ static void performOp(
logger->startWork();
{
FramedSource source(conn.from);
store->addMultipleToStore(source, RepairFlag{repair}, dontCheckSigs ? NoCheckSigs : CheckSigs);
auto expected = readNum<uint64_t>(source);
for (uint64_t i = 0; i < expected; ++i) {
auto info = WorkerProto::Serialise<ValidPathInfo>::read(
*store,
WorkerProto::ReadConn{
.from = source,
.version = conn.protoVersion.features.contains(WorkerProto::featureVersionedAddToStoreMultiple)
? conn.protoVersion
: WorkerProto::Version{.number = {.major = 1, .minor = 16}},
});
Comment thread
edolstra marked this conversation as resolved.
info.ultimate = false;
EnsureRead wrapper{source, info.narSize};
store->addToStore(info, wrapper, RepairFlag{repair}, dontCheckSigs ? NoCheckSigs : CheckSigs);
wrapper.finish();
}
}
logger->stopWork();
break;
Expand Down
2 changes: 0 additions & 2 deletions src/libstore/include/nix/store/remote-store.hh
Original file line number Diff line number Diff line change
Expand Up @@ -103,8 +103,6 @@ struct RemoteStore : public virtual Store,

void addToStore(const ValidPathInfo & info, Source & nar, RepairFlag repair, CheckSigsFlag checkSigs) override;

void addMultipleToStore(Source & source, RepairFlag repair, CheckSigsFlag checkSigs) override;

void
addMultipleToStore(PathsSource && pathsToCopy, Activity & act, RepairFlag repair, CheckSigsFlag checkSigs) override;

Expand Down
2 changes: 0 additions & 2 deletions src/libstore/include/nix/store/store-api.hh
Original file line number Diff line number Diff line change
Expand Up @@ -571,8 +571,6 @@ public:
/**
* Import multiple paths into the store.
*/
virtual void addMultipleToStore(Source & source, RepairFlag repair = NoRepair, CheckSigsFlag checkSigs = CheckSigs);

virtual void addMultipleToStore(
PathsSource && pathsToCopy, Activity & act, RepairFlag repair = NoRepair, CheckSigsFlag checkSigs = CheckSigs);

Expand Down
2 changes: 0 additions & 2 deletions src/libstore/include/nix/store/worker-protocol-connection.hh
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,6 @@ struct WorkerProto::BasicConnection
return WorkerProto::ReadConn{
.from = from,
.version = protoVersion,
.provenance = protoVersion.features.contains(WorkerProto::featureProvenance),
};
}

Expand All @@ -53,7 +52,6 @@ struct WorkerProto::BasicConnection
return WorkerProto::WriteConn{
.to = to,
.version = protoVersion,
.provenance = protoVersion.features.contains(WorkerProto::featureProvenance),
};
}
};
Expand Down
3 changes: 1 addition & 2 deletions src/libstore/include/nix/store/worker-protocol.hh
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,7 @@ struct WorkerProto

static constexpr std::string_view featureQueryActiveBuilds = "queryActiveBuilds";
static constexpr std::string_view featureProvenance = "provenance";
static constexpr std::string_view featureVersionedAddToStoreMultiple = "versionedAddToStoreMultiple";

/**
* A unidirectional read connection, to be used by the read half of the
Expand All @@ -121,7 +122,6 @@ struct WorkerProto
Source & from;
const Version & version;
bool shortStorePaths = false;
bool provenance = false;
};

/**
Expand All @@ -133,7 +133,6 @@ struct WorkerProto
Sink & to;
const Version & version;
bool shortStorePaths = false;
bool provenance = false;
};

/**
Expand Down
24 changes: 12 additions & 12 deletions src/libstore/remote-store.cc
Original file line number Diff line number Diff line change
Expand Up @@ -469,6 +469,13 @@ void RemoteStore::addToStore(const ValidPathInfo & info, Source & source, Repair
void RemoteStore::addMultipleToStore(
PathsSource && pathsToCopy, Activity & act, RepairFlag repair, CheckSigsFlag checkSigs)
{
if (getConnection()->protoVersion.number < WorkerProto::Version::Number{1, 32}) {
Store::addMultipleToStore(std::move(pathsToCopy), act, repair, checkSigs);
return;
}

auto conn(getConnection());

// `addMultipleToStore` is single threaded
size_t bytesExpected = 0;
for (auto & [pathInfo, _] : pathsToCopy) {
Expand All @@ -489,25 +496,18 @@ void RemoteStore::addMultipleToStore(
*this,
WorkerProto::WriteConn{
.to = sink,
.version = {.number = {.major = 1, .minor = 16}},
.version = conn->protoVersion.features.contains(WorkerProto::featureVersionedAddToStoreMultiple)
? conn->protoVersion
: WorkerProto::Version{.number = {.major = 1, .minor = 16}},
},
pathInfo);
pathSource->drainInto(sink);
pathsToCopy.pop_back();
}
});

addMultipleToStore(*source, repair, checkSigs);
}

void RemoteStore::addMultipleToStore(Source & source, RepairFlag repair, CheckSigsFlag checkSigs)
{
if (getConnection()->protoVersion >= WorkerProto::Version{.number = {1, 32}}) {
auto conn(getConnection());
conn->to << WorkerProto::Op::AddMultipleToStore << repair << !checkSigs;
conn.withFramedSink([&](Sink & sink) { source.drainInto(sink); });
} else
Store::addMultipleToStore(source, repair, checkSigs);
conn->to << WorkerProto::Op::AddMultipleToStore << repair << !checkSigs;
conn.withFramedSink([&](Sink & sink) { source->drainInto(sink); });
}

void RemoteStore::registerDrvOutput(const Realisation & info)
Expand Down
22 changes: 0 additions & 22 deletions src/libstore/store-api.cc
Original file line number Diff line number Diff line change
Expand Up @@ -14,9 +14,6 @@
#include "nix/util/callback.hh"
#include "nix/util/git.hh"
#include "nix/util/posix-source-accessor.hh"
// FIXME this should not be here, see TODO below on
// `addMultipleToStore`.
#include "nix/store/worker-protocol.hh"
#include "nix/util/signals.hh"
#include "nix/util/environment-variables.hh"
#include "nix/util/file-system.hh"
Expand Down Expand Up @@ -224,25 +221,6 @@ void Store::addMultipleToStore(PathsSource && pathsToCopy, Activity & act, Repai
});
}

void Store::addMultipleToStore(Source & source, RepairFlag repair, CheckSigsFlag checkSigs)
{
auto expected = readNum<uint64_t>(source);
for (uint64_t i = 0; i < expected; ++i) {
// FIXME we should not be using the worker protocol here, let
// alone the worker protocol with a hard-coded version!
auto info = WorkerProto::Serialise<ValidPathInfo>::read(
*this,
WorkerProto::ReadConn{
.from = source,
.version = {.number = {.major = 1, .minor = 16}},
});
info.ultimate = false;
EnsureRead wrapper{source, info.narSize};
addToStore(info, wrapper, repair, checkSigs);
wrapper.finish();
}
}

/*
The aim of this function is to compute in one pass the correct ValidPathInfo for
the files that we are trying to add to the store. To accomplish that in one
Expand Down
5 changes: 3 additions & 2 deletions src/libstore/worker-protocol.cc
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ const WorkerProto::Version WorkerProto::latest = {
{
std::string{WorkerProto::featureQueryActiveBuilds},
std::string{WorkerProto::featureProvenance},
std::string{WorkerProto::featureVersionedAddToStoreMultiple},
},
};

Expand Down Expand Up @@ -352,7 +353,7 @@ UnkeyedValidPathInfo WorkerProto::Serialise<UnkeyedValidPathInfo>::read(const St
info.sigs = WorkerProto::Serialise<std::set<Signature>>::read(store, conn);
info.ca = ContentAddress::parseOpt(readString(conn.from));
}
if (conn.provenance)
if (conn.version.features.contains(WorkerProto::featureProvenance))
info.provenance = Provenance::from_json_str_optional(readString(conn.from));
return info;
}
Expand All @@ -369,7 +370,7 @@ void WorkerProto::Serialise<UnkeyedValidPathInfo>::write(
WorkerProto::write(store, conn, pathInfo.sigs);
conn.to << renderContentAddress(pathInfo.ca);
}
if (conn.provenance)
if (conn.version.features.contains(WorkerProto::featureProvenance))
conn.to << (pathInfo.provenance ? pathInfo.provenance->to_json_str() : "");
}

Expand Down
Loading