Factor out parallel processing of work items that have dependencies

This commit is contained in:
Eelco Dolstra 2016-04-22 20:50:06 +02:00
parent 91539d305f
commit c879a20850
2 changed files with 78 additions and 55 deletions

View file

@ -58,70 +58,33 @@ struct CmdCopy : StorePathsCommand
progressBar.updateStatus(showProgress());
struct Graph
{
std::set<Path> left;
std::map<Path, std::set<Path>> refs, rrefs;
};
Sync<Graph> graph_;
{
auto graph(graph_.lock());
graph->left = PathSet(storePaths.begin(), storePaths.end());
}
ThreadPool pool;
std::function<void(const Path &)> doPath;
processGraph<Path>(pool,
PathSet(storePaths.begin(), storePaths.end()),
doPath = [&](const Path & storePath) {
checkInterrupt();
[&](const Path & storePath) {
return srcStore->queryPathInfo(storePath)->references;
},
if (!dstStore->isValidPath(storePath)) {
auto activity(progressBar.startActivity(format("copying %s...") % storePath));
[&](const Path & storePath) {
checkInterrupt();
StringSink sink;
srcStore->exportPaths({storePath}, false, sink);
if (!dstStore->isValidPath(storePath)) {
auto activity(progressBar.startActivity(format("copying %s...") % storePath));
StringSource source(*sink.s);
dstStore->importPaths(false, source, 0);
StringSink sink;
srcStore->exportPaths({storePath}, false, sink);
done++;
} else
total--;
StringSource source(*sink.s);
dstStore->importPaths(false, source, 0);
progressBar.updateStatus(showProgress());
done++;
} else
total--;
/* Enqueue all paths that were waiting for this one. */
{
auto graph(graph_.lock());
graph->left.erase(storePath);
for (auto & rref : graph->rrefs[storePath]) {
auto & refs(graph->refs[rref]);
auto i = refs.find(storePath);
assert(i != refs.end());
refs.erase(i);
if (refs.empty())
pool.enqueue(std::bind(doPath, rref));
}
}
};
/* Build the dependency graph; enqueue all paths with no
dependencies. */
for (auto & storePath : storePaths) {
auto info = srcStore->queryPathInfo(storePath);
{
auto graph(graph_.lock());
for (auto & ref : info->references)
if (ref != storePath && graph->left.count(ref)) {
graph->refs[storePath].insert(ref);
graph->rrefs[ref].insert(storePath);
}
if (graph->refs[storePath].empty())
pool.enqueue(std::bind(doPath, storePath));
}
}
progressBar.updateStatus(showProgress());
});
pool.process();