-
Notifications
You must be signed in to change notification settings - Fork 6
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Added MultiPart Upload to Transaction
- Loading branch information
Showing
9 changed files
with
416 additions
and
57 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Large diffs are not rendered by default.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,178 @@ | ||
#include "catch2/single_include/catch2/catch.hpp" | ||
#include "cloud/minio.hpp" | ||
#include "cloud/provider.hpp" | ||
#include "network/tasked_send_receiver.hpp" | ||
#include "network/transaction.hpp" | ||
#include <cstdlib> | ||
#include <cstring> | ||
#include <filesystem> | ||
#include <fstream> | ||
#include <future> | ||
//--------------------------------------------------------------------------- | ||
// AnyBlob - Universal Cloud Object Storage Library | ||
// Dominik Durner, 2022 | ||
// | ||
// This Source Code Form is subject to the terms of the Mozilla Public License, v. 2.0. | ||
// If a copy of the MPL was not distributed with this file, You can obtain one at http://mozilla.org/MPL/2.0/. | ||
// SPDX-License-Identifier: MPL-2.0 | ||
//--------------------------------------------------------------------------- | ||
namespace anyblob { | ||
namespace test { | ||
//--------------------------------------------------------------------------- | ||
using namespace std; | ||
//--------------------------------------------------------------------------- | ||
TEST_CASE("MinIO Asynchronous Integration") { | ||
// Get the enviornment | ||
const char* bucket = getenv("AWS_S3_BUCKET"); | ||
const char* region = getenv("AWS_S3_REGION"); | ||
const char* endpoint = getenv("AWS_S3_ENDPOINT"); | ||
const char* key = getenv("AWS_S3_ACCESS_KEY"); | ||
const char* secret = getenv("AWS_S3_SECRET_ACCESS_KEY"); | ||
|
||
REQUIRE(bucket); | ||
REQUIRE(region); | ||
REQUIRE(endpoint); | ||
REQUIRE(key); | ||
REQUIRE(secret); | ||
|
||
auto stringGen = [](auto len) { | ||
static constexpr auto chars = "0123456789ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz"; | ||
auto resultString = string(len, '\0'); | ||
generate_n(begin(resultString), len, [&]() { return chars[rand() % strlen(chars)]; }); | ||
return resultString; | ||
}; | ||
|
||
// The file to be uploaded and downloaded | ||
string bucketName = "minio://"; | ||
bucketName = bucketName + endpoint + "/" + bucket + ":" + region; | ||
string fileName[]{"test.txt", "long.txt"}; | ||
string longText = stringGen(1 << 24); | ||
string content[]{"Hello World!", longText}; | ||
|
||
// Create a new task group (18 concurrent request maximum, and up to 1024 outstanding submissions) | ||
anyblob::network::TaskedSendReceiverGroup group(18, 1024); | ||
|
||
// Create an AnyBlob scheduler object for the group | ||
anyblob::network::TaskedSendReceiver sendReceiver(group); | ||
|
||
// Async thread | ||
future<void> asyncSendReceiverThread; | ||
|
||
auto runLambda = [&]() { | ||
// Runs the download thread to asynchronously retrieve data | ||
sendReceiver.run(); | ||
}; | ||
asyncSendReceiverThread = async(launch::async, runLambda); | ||
|
||
// Create the provider for the corresponding filename | ||
auto provider = anyblob::cloud::Provider::makeProvider(bucketName, false, key, secret, &sendReceiver); | ||
{ | ||
// Check the upload for success | ||
std::atomic<uint16_t> finishedMessages = 0; | ||
auto checkSuccess = [&finishedMessages](anyblob::network::MessageResult& result) { | ||
// Sucessful request | ||
REQUIRE(result.success()); | ||
finishedMessages++; | ||
}; | ||
|
||
// Create the put request | ||
anyblob::network::Transaction putTxn(provider.get()); | ||
for (auto i = 0u; i < 2; i++) | ||
putTxn.putObjectRequest(checkSuccess, fileName[i], content[i].data(), content[i].size()); | ||
|
||
// Upload the request asynchronously | ||
putTxn.processAsync(group); | ||
|
||
// Wait for the upload | ||
while (finishedMessages != 2) | ||
usleep(100); | ||
} | ||
{ | ||
// Check the upload for success | ||
std::atomic<uint16_t> finishedMessages = 0; | ||
auto checkSuccess = [&finishedMessages](anyblob::network::MessageResult& result) { | ||
// Sucessful request | ||
REQUIRE(result.success()); | ||
finishedMessages++; | ||
}; | ||
|
||
// Create the multipart put request | ||
auto minio = static_cast<anyblob::cloud::MinIO*>(provider.get()); | ||
minio->setMultipartUploadSize(6ull << 20); | ||
anyblob::network::Transaction putTxn(provider.get()); | ||
for (auto i = 0u; i < 2; i++) | ||
putTxn.putObjectRequest(checkSuccess, fileName[i], content[i].data(), content[i].size()); | ||
|
||
while (finishedMessages != 2) { | ||
// Upload the new request asynchronously | ||
putTxn.processAsync(group); | ||
// Wait for the upload | ||
usleep(100); | ||
} | ||
} | ||
{ | ||
std::atomic<uint16_t> finishedMessages = 0; | ||
// Create the get request | ||
anyblob::network::Transaction getTxn(provider.get()); | ||
for (auto i = 0u; i < 2; i++) { | ||
// Check the download for success | ||
auto checkSuccess = [&finishedMessages, &content, i](anyblob::network::MessageResult& result) { | ||
// Sucessful request | ||
REQUIRE(result.success()); | ||
// Simple string_view interface | ||
REQUIRE(!content[i].compare(result.getResult())); | ||
// Check from the other side too | ||
REQUIRE(!result.getResult().compare(content[i])); | ||
// Check the size | ||
REQUIRE(result.getSize() == content[i].size()); | ||
|
||
// Advanced raw interface | ||
// Note that the data lies in the data buffer but after the offset to skip the HTTP header | ||
// Note that the size is already without the header, so the full request has size + offset length | ||
string_view rawDataString(reinterpret_cast<const char*>(result.getData()) + result.getOffset(), result.getSize()); | ||
REQUIRE(!content[i].compare(rawDataString)); | ||
REQUIRE(!rawDataString.compare(result.getResult())); | ||
REQUIRE(!rawDataString.compare(content[i])); | ||
finishedMessages++; | ||
}; | ||
|
||
getTxn.getObjectRequest(std::move(checkSuccess), fileName[i]); | ||
} | ||
|
||
// Retrieve the request asynchronously | ||
getTxn.processAsync(group); | ||
|
||
// Wait for the download | ||
while (finishedMessages != 2) | ||
usleep(100); | ||
} | ||
{ | ||
// Check the delete for success | ||
std::atomic<uint16_t> finishedMessages = 0; | ||
auto checkSuccess = [&finishedMessages](anyblob::network::MessageResult& result) { | ||
// Sucessful request | ||
REQUIRE(result.success()); | ||
finishedMessages++; | ||
}; | ||
|
||
// Create the delete request | ||
anyblob::network::Transaction deleteTxn(provider.get()); | ||
for (auto i = 0u; i < 2; i++) | ||
deleteTxn.deleteObjectRequest(checkSuccess, fileName[i]); | ||
|
||
// Process the request asynchronously | ||
deleteTxn.processAsync(group); | ||
|
||
// Wait for the deletion | ||
while (finishedMessages != 2) | ||
usleep(100); | ||
} | ||
|
||
// Stop the send receiver daemon | ||
sendReceiver.stop(); | ||
// Join the thread | ||
asyncSendReceiverThread.get(); | ||
} | ||
//--------------------------------------------------------------------------- | ||
} // namespace test | ||
} // namespace anyblob |
Oops, something went wrong.