mirror of
https://github.com/revng/revng
synced 2026-06-21 14:07:57 +00:00
S3StorageClient: implement parallel upload
Allow uploads to S3 to be executed in parallel, this should reduce the time it takes for a `save` operation to conclude.
This commit is contained in:
committed by
Alessandro Di Federico
parent
761e84fc9f
commit
a06b380ce0
@@ -24,4 +24,17 @@ inline void cantFail(std::error_code EC) {
|
||||
revng_assert(not EC);
|
||||
}
|
||||
|
||||
template<std::ranges::range T>
|
||||
inline llvm::Error joinErrors(T &Container) {
|
||||
auto Iter = Container.begin();
|
||||
llvm::Error Result{ std::move(*Iter) };
|
||||
if (std::distance(Container.begin(), Container.end()) == 1) {
|
||||
return Result;
|
||||
}
|
||||
for (Iter++; Iter < Container.end(); Iter++) {
|
||||
Result = llvm::joinErrors(std::move(Result), std::move(*Iter));
|
||||
}
|
||||
return Result;
|
||||
}
|
||||
|
||||
} // namespace revng
|
||||
|
||||
@@ -71,6 +71,7 @@ public:
|
||||
};
|
||||
|
||||
bool SDKIsInitialized = false;
|
||||
std::optional<llvm::ThreadPool> ThreadPool;
|
||||
|
||||
void initializeSDK() {
|
||||
revng_assert(!SDKIsInitialized);
|
||||
@@ -84,6 +85,8 @@ void initializeSDK() {
|
||||
};
|
||||
Aws::InitAPI(Options);
|
||||
OnQuit->add([Options = std::move(Options)] { Aws::ShutdownAPI(Options); });
|
||||
|
||||
ThreadPool.emplace(llvm::hardware_concurrency(8));
|
||||
}
|
||||
|
||||
} // namespace
|
||||
@@ -156,6 +159,41 @@ public:
|
||||
llvm::MemoryBuffer &buffer() override { return *Buffer; };
|
||||
};
|
||||
|
||||
class AsyncUploadTask {
|
||||
private:
|
||||
Aws::S3::S3Client &Client;
|
||||
Aws::S3::Model::PutObjectRequest Request;
|
||||
std::string Path;
|
||||
std::string NewFilename;
|
||||
TemporaryFile TempFile;
|
||||
|
||||
public:
|
||||
AsyncUploadTask(Aws::S3::S3Client &Client,
|
||||
Aws::S3::Model::PutObjectRequest &&Request,
|
||||
llvm::StringRef Path,
|
||||
llvm::StringRef NewFilename,
|
||||
TemporaryFile &&TempFile) :
|
||||
Client(Client),
|
||||
Request(std::move(Request)),
|
||||
Path(Path),
|
||||
NewFilename(NewFilename),
|
||||
TempFile(std::move(TempFile)) {}
|
||||
|
||||
S3StorageClient::UploadResult operator()() {
|
||||
Aws::S3::Model::PutObjectOutcome Result;
|
||||
auto File = std::make_shared<Aws::FStream>(TempFile.path().str(),
|
||||
std::ios_base::in
|
||||
| std::ios_base::binary);
|
||||
if (File->fail()) {
|
||||
return { Result, Path, NewFilename, "Could not open temporary file" };
|
||||
}
|
||||
|
||||
Request.SetBody(File);
|
||||
Result = Client.PutObject(Request);
|
||||
return { Result, Path, NewFilename, std::string{} };
|
||||
}
|
||||
};
|
||||
|
||||
class S3WritableFile : public WritableFile {
|
||||
private:
|
||||
TemporaryFile TempFile;
|
||||
@@ -189,19 +227,15 @@ public:
|
||||
std::string NewFilename = generateNewFilename(Path);
|
||||
Request.SetKey(Client.resolvePath(NewFilename));
|
||||
|
||||
auto File = std::make_shared<Aws::FStream>(TempFile.path().str(),
|
||||
std::ios_base::in
|
||||
| std::ios_base::binary);
|
||||
if (File->fail()) {
|
||||
return revng::createError("Could not open temporary file");
|
||||
}
|
||||
|
||||
Request.SetBody(File);
|
||||
Aws::S3::Model::PutObjectOutcome Result = Client.Client.PutObject(Request);
|
||||
if (not Result.IsSuccess())
|
||||
return toError(Result);
|
||||
|
||||
Client.FilenameMap[Path] = NewFilename;
|
||||
auto Task = std::make_shared<AsyncUploadTask>(Client.Client,
|
||||
std::move(Request),
|
||||
Path,
|
||||
NewFilename,
|
||||
std::move(TempFile));
|
||||
using UploadResult = S3StorageClient::UploadResult;
|
||||
std::shared_future<UploadResult>
|
||||
ResultFuture = Client.TaskGroup.async([Task]() { return (*Task)(); });
|
||||
Client.PendingUploads.push_back(std::move(ResultFuture));
|
||||
return llvm::Error::success();
|
||||
}
|
||||
};
|
||||
@@ -217,7 +251,8 @@ public:
|
||||
Aws::Auth::AWSCredentials GetAWSCredentials() override { return Credentials; }
|
||||
};
|
||||
|
||||
S3StorageClient::S3StorageClient(llvm::StringRef RawURL) {
|
||||
S3StorageClient::S3StorageClient(llvm::StringRef RawURL) :
|
||||
TaskGroup(*ThreadPool) {
|
||||
// Url format is:
|
||||
// s3://<username>:<password>@<region>+<host:port>/<bucket name>/<path>
|
||||
revng_assert(isS3URL(RawURL));
|
||||
@@ -426,6 +461,22 @@ S3StorageClient::getWritableFile(llvm::StringRef Path,
|
||||
}
|
||||
|
||||
llvm::Error S3StorageClient::commit() {
|
||||
TaskGroup.wait();
|
||||
|
||||
std::vector<llvm::Error> Errors;
|
||||
for (const auto &Element : PendingUploads) {
|
||||
const UploadResult &Result = Element.get();
|
||||
if (not Result.Error.empty())
|
||||
Errors.push_back(revng::createError(Result.Error));
|
||||
else if (not Result.Outcome.IsSuccess())
|
||||
Errors.push_back(toError(Result.Outcome));
|
||||
else
|
||||
FilenameMap[Result.Path] = Result.NewPath;
|
||||
}
|
||||
if (Errors.size() > 0) {
|
||||
return joinErrors(Errors);
|
||||
}
|
||||
|
||||
std::string SerializedIndex;
|
||||
|
||||
{
|
||||
|
||||
@@ -9,20 +9,32 @@
|
||||
|
||||
#include "llvm/ADT/StringMap.h"
|
||||
#include "llvm/ADT/StringRef.h"
|
||||
#include "llvm/Support/ThreadPool.h"
|
||||
|
||||
#include "revng/Storage/StorageClient.h"
|
||||
|
||||
namespace revng {
|
||||
|
||||
class AsyncUploadTask;
|
||||
class S3WritableFile;
|
||||
|
||||
class S3StorageClient : public StorageClient {
|
||||
private:
|
||||
struct UploadResult {
|
||||
Aws::S3::Model::PutObjectOutcome Outcome;
|
||||
std::string Path;
|
||||
std::string NewPath;
|
||||
std::string Error;
|
||||
};
|
||||
|
||||
private:
|
||||
Aws::Auth::AWSCredentials Credentials;
|
||||
Aws::S3::S3Client Client;
|
||||
std::string Bucket;
|
||||
std::string SubPath;
|
||||
std::string RedactedURL;
|
||||
llvm::ThreadPoolTaskGroup TaskGroup;
|
||||
std::vector<std::shared_future<UploadResult>> PendingUploads;
|
||||
llvm::StringMap<std::string> FilenameMap;
|
||||
static constexpr auto IndexName = "index.yml";
|
||||
|
||||
@@ -61,6 +73,7 @@ private:
|
||||
std::string dumpString() const override;
|
||||
std::string resolvePath(llvm::StringRef Path);
|
||||
friend class S3WritableFile;
|
||||
friend class AsyncUploadTask;
|
||||
};
|
||||
|
||||
} // namespace revng
|
||||
|
||||
Reference in New Issue
Block a user