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
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -109,3 +109,6 @@ third-party/folly/
output

.pre-commit-config.yaml

# Purger tests read MANIFEST-<epoch> objects into the working directory
/MANIFEST-*
8 changes: 6 additions & 2 deletions CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -976,11 +976,12 @@ set(SOURCES
cloud/cloud_file_system_impl.cc
cloud/cloud_log_controller.cc
cloud/manifest_reader.cc
# cloud/purge.cc
cloud/improved_purger.cc
cloud/eloq_purger.cc
cloud/file_number_guard.cc
cloud/cloud_manifest.cc
cloud/cloud_scheduler.cc
cloud/cloud_storage_provider.cc
cloud/cloud_file_cache.cc
cloud/cloud_file_deletion_scheduler.cc
db/db_impl/replication_codec.cc)

Expand Down Expand Up @@ -1318,6 +1319,9 @@ if(WITH_TESTS)
cloud/db_cloud_test.cc
cloud/cloud_manifest_test.cc
cloud/cloud_scheduler_test.cc
cloud/eloq_purger_test.cc
cloud/eloq_purger_integration_test.cc
cloud/file_number_guard_test.cc
cloud/replication_test.cc
cache/tiered_secondary_cache_test.cc
db/blob/blob_counting_iterator_test.cc
Expand Down
9 changes: 9 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -1901,6 +1901,15 @@ cloud_manifest_test: cloud/cloud_manifest_test.o $(TEST_LIBRARY) $(LIBRARY)
cloud_scheduler_test: cloud/cloud_scheduler_test.o $(TEST_LIBRARY) $(LIBRARY)
$(AM_LINK)

eloq_purger_test: cloud/eloq_purger_test.o $(TEST_LIBRARY) $(LIBRARY)
$(AM_LINK)

file_number_guard_test: cloud/file_number_guard_test.o $(TEST_LIBRARY) $(LIBRARY)
$(AM_LINK)

eloq_purger_integration_test: cloud/eloq_purger_integration_test.o $(TEST_LIBRARY) $(LIBRARY)
$(AM_LINK)

iostats_context_test: $(OBJ_DIR)/monitoring/iostats_context_test.o $(TEST_LIBRARY) $(LIBRARY)
$(AM_V_CCLD)$(CXX) $^ $(EXEC_LDFLAGS) -o $@ $(LDFLAGS)

Expand Down
84 changes: 84 additions & 0 deletions cloud/aws/aws_s3.cc
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,13 @@
#include <aws/s3/model/CreateBucketConfiguration.h>
#include <aws/s3/model/CreateBucketRequest.h>
#include <aws/s3/model/CreateBucketResult.h>
#include <aws/s3/model/Delete.h>
#include <aws/s3/model/DeleteBucketRequest.h>
#include <aws/s3/model/DeleteObjectRequest.h>
#include <aws/s3/model/DeleteObjectResult.h>
#include <aws/s3/model/DeleteObjectsRequest.h>
#include <aws/s3/model/DeleteObjectsResult.h>
#include <aws/s3/model/ObjectIdentifier.h>
#include <aws/s3/model/GetBucketVersioningRequest.h>
#include <aws/s3/model/GetBucketVersioningResult.h>
#include <aws/s3/model/GetObjectRequest.h>
Expand Down Expand Up @@ -184,6 +188,15 @@ class AwsS3ClientWrapper {
return outcome;
}

Aws::S3::Model::DeleteObjectsOutcome DeleteCloudObjects(
const Aws::S3::Model::DeleteObjectsRequest& request) {
CloudRequestCallbackGuard t(cloud_request_callback_.get(),
CloudRequestOpType::kDeleteOp);
auto outcome = client_->DeleteObjects(request);
t.SetSuccess(outcome.IsSuccess());
return outcome;
}

Aws::S3::Model::CopyObjectOutcome CopyCloudObject(
const Aws::S3::Model::CopyObjectRequest& request) {
CloudRequestCallbackGuard t(cloud_request_callback_.get(),
Expand Down Expand Up @@ -392,6 +405,10 @@ class S3StorageProvider : public CloudStorageProviderImpl {
const std::string& object_path) override;
IOStatus DeleteCloudObject(const std::string& bucket_name,
const std::string& object_path) override;
IOStatus DeleteCloudObjects(const std::string& bucket_name,
const std::vector<std::string>& object_paths,
size_t* deleted_count,
size_t* failed_count) override;
IOStatus ListCloudObjects(const std::string& bucket_name,
const std::string& object_path,
std::vector<std::string>* result) override;
Expand Down Expand Up @@ -634,6 +651,73 @@ IOStatus S3StorageProvider::DeleteCloudObject(const std::string& bucket_name,
return st;
}

IOStatus S3StorageProvider::DeleteCloudObjects(
const std::string& bucket_name,
const std::vector<std::string>& object_paths, size_t* deleted_count,
size_t* failed_count) {
assert(deleted_count != nullptr && failed_count != nullptr);
IOStatus first_error;
// S3 DeleteObjects accepts at most 1000 keys per request. A key that does
// not exist is reported as deleted, matching the interface contract.
constexpr size_t kMaxKeysPerBatch = 1000;

for (size_t begin = 0; begin < object_paths.size();
begin += kMaxKeysPerBatch) {
size_t end = begin + kMaxKeysPerBatch;
if (end > object_paths.size()) {
end = object_paths.size();
}
size_t batch_size = end - begin;

Aws::S3::Model::Delete to_delete;
for (size_t i = begin; i < end; ++i) {
to_delete.AddObjects(
Aws::S3::Model::ObjectIdentifier().WithKey(
ToAwsString(object_paths[i])));
}
// Quiet mode: the response lists only the keys that failed.
to_delete.SetQuiet(true);

Aws::S3::Model::DeleteObjectsRequest request;
request.SetBucket(ToAwsString(bucket_name));
request.SetDelete(std::move(to_delete));

auto outcome = s3client_->DeleteCloudObjects(request);
if (!outcome.IsSuccess()) {
const Aws::Client::AWSError<Aws::S3::S3Errors>& error =
outcome.GetError();
std::string errmsg(error.GetMessage().c_str());
Log(InfoLogLevel::ERROR_LEVEL, cfs_->GetLogger(),
"[s3] DeleteObjects batch of %zu keys failed in bucket %s: %s",
batch_size, bucket_name.c_str(), errmsg.c_str());
*failed_count += batch_size;
if (first_error.ok()) {
first_error = IOStatus::IOError(bucket_name, errmsg.c_str());
}
continue;
}

const auto& errors = outcome.GetResult().GetErrors();
for (const auto& error : errors) {
Log(InfoLogLevel::ERROR_LEVEL, cfs_->GetLogger(),
"[s3] DeleteObjects failed to delete %s/%s: %s %s",
bucket_name.c_str(), error.GetKey().c_str(),
error.GetCode().c_str(), error.GetMessage().c_str());
if (first_error.ok()) {
first_error = IOStatus::IOError(std::string(error.GetKey().c_str()),
std::string(error.GetMessage().c_str()));
}
}
*failed_count += errors.size();
*deleted_count += batch_size - errors.size();

Log(InfoLogLevel::INFO_LEVEL, cfs_->GetLogger(),
"[s3] DeleteObjects deleted %zu of %zu keys from bucket %s",
batch_size - errors.size(), batch_size, bucket_name.c_str());
}
return first_error;
}

//
// Appends the names of all children of the specified path from S3
// into the result set.
Expand Down
40 changes: 40 additions & 0 deletions cloud/cloud_file_system.cc
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,12 @@ void CloudFileSystemOptions::Dump(Logger* log) const {
create_bucket_if_missing ? "true" : "false");
Header(log, " COptions.run_purger: %s",
run_purger ? "true" : "false");
Header(log, " COptions.publish_file_number_guard: %s",
publish_file_number_guard ? "true" : "false");
Header(log, " COptions.guard_publish_interval: %llu ms",
static_cast<unsigned long long>(guard_publish_interval.count()));
Header(log, " COptions.guard_entry_duration: %llu ms",
static_cast<unsigned long long>(guard_entry_duration.count()));
Header(log, " COptions.resync_on_open: %s",
resync_on_open ? "true" : "false");
Header(log, " COptions.skip_dbid_verification: %s",
Expand Down Expand Up @@ -302,6 +308,31 @@ int offset_of(T1 CloudFileSystemOptions::*member) {
return int(size_t(&(dummy_ceo_options.*member)) - size_t(&dummy_ceo_options));
}

static OptionTypeInfo MillisecondsOption(int offset) {
return {offset,
OptionType::kInt64T,
OptionVerificationType::kNormal,
OptionTypeFlags::kNone,
[](const ConfigOptions& /*opts*/, const std::string& /*name*/,
const std::string& value, void* addr) {
*static_cast<std::chrono::milliseconds*>(addr) =
std::chrono::milliseconds(ParseInt64(value));
return Status::OK();
},
[](const ConfigOptions& /*opts*/, const std::string& /*name*/,
const void* addr, std::string* value) {
*value = std::to_string(
static_cast<const std::chrono::milliseconds*>(addr)->count());
return Status::OK();
},
[](const ConfigOptions& /*opts*/, const std::string& /*name*/,
const void* addr1, const void* addr2,
std::string* /*mismatch*/) {
return *static_cast<const std::chrono::milliseconds*>(addr1) ==
*static_cast<const std::chrono::milliseconds*>(addr2);
}};
}

const std::unordered_map<std::string, OptionTypeInfo>
CloudFileSystemOptions::cloud_fs_option_type_info = {
{"keep_local_sst_files",
Expand Down Expand Up @@ -335,6 +366,15 @@ const std::unordered_map<std::string, OptionTypeInfo>
{"purger_periodicity_ms",
{offset_of(&CloudFileSystemOptions::purger_periodicity_millis),
OptionType::kUInt64T}},
{"publish_file_number_guard",
Comment thread
liunyl marked this conversation as resolved.
{offset_of(&CloudFileSystemOptions::publish_file_number_guard),
OptionType::kBoolean}},
{"guard_publish_interval_ms",
MillisecondsOption(
offset_of(&CloudFileSystemOptions::guard_publish_interval))},
{"guard_entry_duration_ms",
MillisecondsOption(
offset_of(&CloudFileSystemOptions::guard_entry_duration))},

{"provider",
{offset_of(&CloudFileSystemOptions::storage_provider),
Expand Down
9 changes: 9 additions & 0 deletions cloud/cloud_file_system_impl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,9 @@ CloudFileSystemImpl::~CloudFileSystemImpl() {
if (cloud_fs_options.cloud_log_controller) {
cloud_fs_options.cloud_log_controller->StopTailingStream();
}
// Cancel the file number guard's scheduled publish before the members its
// callback reaches through (storage provider, cloud manifest) go away.
StopFileNumberGuard();
StopPurger();
FileCachePurge();
cloud_fs_options.cloud_log_controller.reset();
Expand Down Expand Up @@ -2343,6 +2346,12 @@ Status CloudFileSystemImpl::CheckValidity() const {
cloud_fs_options.dest_bucket.GetObjectPath().empty()) {
return Status::InvalidArgument(
"Must specify both dest bucket name and path");
} else if (cloud_fs_options.guard_publish_interval.count() <= 0) {
return Status::InvalidArgument(
"guard_publish_interval_ms must be greater than zero");
} else if (cloud_fs_options.guard_entry_duration.count() <= 0) {
return Status::InvalidArgument(
"guard_entry_duration_ms must be greater than zero");
} else if (!cloud_fs_options.storage_provider) {
return Status::InvalidArgument(
"Cloud environment requires a storage provider");
Expand Down
54 changes: 52 additions & 2 deletions cloud/cloud_storage_provider.cc
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@

#include <cinttypes>

#include "cloud/file_number_guard.h"
#include "cloud/filename.h"
#include "file/filename.h"
#include "rocksdb/cloud/cloud_file_system.h"
Expand Down Expand Up @@ -113,7 +114,6 @@ CloudStorageWritableFileImpl::CloudStorageWritableFileImpl(
auto fname_no_epoch = RemoveEpoch(fname_);
// Is this a manifest file?
is_manifest_ = IsManifestFile(fname_no_epoch);
assert(IsSstFile(fname_no_epoch) || is_manifest_);

Log(InfoLogLevel::DEBUG_LEVEL, cfs_->GetLogger(),
"[%s] CloudWritableFile bucket %s opened local file %s "
Expand Down Expand Up @@ -176,7 +176,37 @@ IOStatus CloudStorageWritableFileImpl::Close(const IOOptions& opts,
local_file_.reset();

if (!is_manifest_) {
status_ = cfs_->CopyLocalFileToDest(fname_, cloud_fname_);
auto* cfs_impl = dynamic_cast<CloudFileSystemImpl*>(cfs_);
auto publisher =
cfs_impl == nullptr ? nullptr : cfs_impl->GetFileNumberGuardPublisher();
if (publisher && !publisher->CoversFile(fname_)) {
// Written outside the DB's own SST directories (a checkpoint copy, a
// column family export): not a DB-visible SST. No flush/compaction job
// registers it and the purger never evaluates it as this DB's live or
// obsolete file, so the guard neither protects nor blocks it.
Log(InfoLogLevel::DEBUG_LEVEL, cfs_->GetLogger(),
"[%s] CloudWritableFile %s is outside the DB directories; skipping "
"file number guard",
Name(), fname_.c_str());
publisher.reset();
}
uint64_t file_number = 0;
FileType file_type;
const std::string logical_name = basename(RemoveEpoch(fname_));
const bool parsed = ParseFileName(logical_name, &file_number, &file_type);
if (publisher && !parsed && IsSstFile(logical_name)) {
status_ = IOStatus::InvalidArgument("cannot parse SST file number",
logical_name);
return status_;
}
if (publisher && parsed && file_type == kTableFile) {
Status protection = publisher->ProtectFileUpload(file_number, [&] {
return cfs_->CopyLocalFileToDest(fname_, cloud_fname_);
});
status_ = status_to_io_status(std::move(protection));
} else {
status_ = cfs_->CopyLocalFileToDest(fname_, cloud_fname_);
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
if (!status_.ok()) {
Log(InfoLogLevel::ERROR_LEVEL, cfs_->GetLogger(),
"[%s] CloudWritableFile closing PutObject failed on local file %s",
Expand Down Expand Up @@ -266,6 +296,26 @@ IOStatus CloudStorageWritableFileImpl::Sync(const IOOptions& opts,

CloudStorageProvider::~CloudStorageProvider() {}

IOStatus CloudStorageProvider::DeleteCloudObjects(
const std::string& bucket_name,
const std::vector<std::string>& object_paths, size_t* deleted_count,
size_t* failed_count) {
assert(deleted_count != nullptr && failed_count != nullptr);
IOStatus first_error;
for (const auto& object_path : object_paths) {
auto st = DeleteCloudObject(bucket_name, object_path);
if (st.ok() || st.IsNotFound()) {
++(*deleted_count);
} else {
++(*failed_count);
if (first_error.ok()) {
first_error = st;
}
}
}
return first_error;
}

Status CloudStorageProvider::CreateFromString(
const ConfigOptions& /*config_options*/, const std::string& id,
std::shared_ptr<CloudStorageProvider>* provider) {
Expand Down
Loading