github-actions[bot] commented on code in PR #66350:
URL: https://github.com/apache/doris/pull/66350#discussion_r3712258420
##########
cloud/test/recycler_batch_delete_test.cpp:
##########
@@ -17,387 +17,193 @@
#include <gtest/gtest.h>
-#include <atomic>
#include <memory>
-#include <optional>
#include <string>
+#include <utility>
#include <vector>
-#include "common/config.h"
-#include "common/logging.h"
-#include "common/simple_thread_pool.h"
-#include "recycler/obj_storage_client.h"
+#include "cpp/client/obj_storage_client.h"
-using namespace doris;
+namespace doris {
+namespace {
-namespace doris::cloud {
-
-// Mock ObjectListIterator for testing
-class MockObjectListIterator : public ObjectListIterator {
+class MockObjStorageBackend final : public ObjStorageBackend {
public:
- MockObjectListIterator(std::vector<ObjectMeta> objects, int fail_after =
-1)
- : objects_(std::move(objects)), fail_after_(fail_after) {}
-
- bool is_valid() override { return is_valid_; }
-
- bool has_next() override {
- if (!is_valid_) return false;
- return current_index_ < objects_.size();
+ MockObjStorageBackend(std::vector<ObjectMeta> objects, size_t batch_size,
+ int iterator_fail_after = -1)
+ : objects_(std::move(objects)),
+ batch_size_(batch_size),
+ iterator_fail_after_(iterator_fail_after) {}
+
+ ObjectStorageUploadResponse create_multipart_upload(const
ObjectStoragePathOptions&) override {
+ return {.resp = ObjectStorageResponse::OK()};
}
-
- std::optional<ObjectMeta> next() override {
- if (!is_valid_ || current_index_ >= objects_.size()) {
- return std::nullopt;
- }
-
- // Simulate iterator becoming invalid after certain number of calls
- if (fail_after_ >= 0 && static_cast<int>(current_index_) >=
fail_after_) {
- is_valid_ = false;
- return std::nullopt;
- }
-
- return objects_[current_index_++];
+ ObjectStorageResponse put_object(const ObjectStoragePathOptions&,
std::string_view) override {
+ return ObjectStorageResponse::OK();
}
-
- void set_invalid() { is_valid_ = false; }
-
-private:
- std::vector<ObjectMeta> objects_;
- size_t current_index_ = 0;
- bool is_valid_ = true;
- int fail_after_ = -1; // -1 means never fail
-};
-
-// Mock ObjStorageClient for testing delete_objects_recursively_
-class MockObjStorageClient : public ObjStorageClient {
-public:
- MockObjStorageClient(std::vector<ObjectMeta> objects, int
iterator_fail_after = -1)
- : objects_(std::move(objects)),
iterator_fail_after_(iterator_fail_after) {}
-
- ObjectStorageResponse put_object(ObjectStoragePathRef path,
std::string_view stream) override {
- return {0};
+ ObjectStorageUploadResponse upload_part(const ObjectStoragePathOptions&,
std::string_view,
+ int) override {
+ return {.resp = ObjectStorageResponse::OK()};
}
-
- ObjectStorageResponse head_object(ObjectStoragePathRef path, ObjectMeta*
res) override {
- return {0};
+ ObjectStorageResponse complete_multipart_upload(
+ const ObjectStoragePathOptions&, const
std::vector<ObjectCompleteMultiPart>&) override {
+ return ObjectStorageResponse::OK();
}
-
- std::unique_ptr<ObjectListIterator> list_objects(ObjectStoragePathRef
path) override {
- return std::make_unique<MockObjectListIterator>(objects_,
iterator_fail_after_);
+ ObjectStorageHeadResponse head_object(const ObjectStoragePathOptions&)
override {
+ return {.resp = ObjectStorageResponse::OK()};
}
-
- ObjectStorageResponse delete_objects(const std::string& bucket,
std::vector<std::string> keys,
- ObjClientOptions option) override {
- delete_calls_++;
- total_keys_deleted_ += keys.size();
-
- // Simulate delete failure if configured
- if (fail_delete_after_ >= 0 && delete_calls_ > fail_delete_after_) {
- return {-1, "simulated delete failure"};
+ ObjectStorageResponse get_object(const ObjectStoragePathOptions&, void*,
size_t, size_t,
+ size_t*) override {
+ return ObjectStorageResponse::OK();
+ }
+ ObjectStorageListPage list_objects(const ObjectStoragePathOptions&,
+ std::string_view continuation_token)
override {
+ const size_t index =
+ continuation_token.empty()
+ ? 0
+ :
static_cast<size_t>(std::stoull(std::string(continuation_token)));
+ if (iterator_fail_after_ >= 0 && index >=
static_cast<size_t>(iterator_fail_after_)) {
+ return {.resp = {.status = {TStatusCode::INTERNAL_ERROR,
"simulated list failure"}}};
}
-
- return {0};
+ ObjectStorageListPage page {.resp = ObjectStorageResponse::OK()};
+ if (index < objects_.size()) {
+ page.objects.emplace_back(objects_[index]);
+ page.has_more = index + 1 < objects_.size();
+ if (page.has_more) {
+ page.continuation_token = std::to_string(index + 1);
+ }
+ }
+ return page;
}
-
- ObjectStorageResponse delete_object(ObjectStoragePathRef path) override {
return {0}; }
-
- ObjectStorageResponse delete_objects_recursively(ObjectStoragePathRef path,
- ObjClientOptions option,
- int64_t expiration_time =
0) override {
- return delete_objects_recursively_(path, option, expiration_time,
1000);
+ ObjectStorageResponse delete_objects(const ObjectStoragePathOptions&,
+ std::vector<std::string> keys)
override {
+ ++delete_calls_;
+ deleted_keys_.insert(deleted_keys_.end(), keys.begin(), keys.end());
+ if (fail_delete_) {
+ return {.status = {TStatusCode::INTERNAL_ERROR, "simulated delete
failure"}};
+ }
+ return ObjectStorageResponse::OK();
}
-
- ObjectStorageResponse get_life_cycle(const std::string& bucket,
- int64_t* expiration_days) override {
- return {0};
+ ObjectStorageResponse delete_object(const ObjectStoragePathOptions&)
override {
+ return ObjectStorageResponse::OK();
}
-
- ObjectStorageResponse check_versioning(const std::string& bucket) override
{ return {0}; }
-
- ObjectStorageResponse abort_multipart_upload(ObjectStoragePathRef path,
- const std::string& upload_id)
override {
- return {0};
+ std::string generate_presigned_url(const ObjectStoragePathOptions&,
int64_t) override {
+ return {};
+ }
+ ObjStorageCapabilities capabilities() const override {
+ return {.max_delete_batch = batch_size_};
}
- // Test helper methods
- int get_delete_calls() const { return delete_calls_; }
- size_t get_total_keys_deleted() const { return total_keys_deleted_; }
- void set_fail_delete_after(int n) { fail_delete_after_ = n; }
+ int delete_calls() const { return delete_calls_; }
+ const std::vector<std::string>& deleted_keys() const { return
deleted_keys_; }
+ void fail_delete() { fail_delete_ = true; }
private:
std::vector<ObjectMeta> objects_;
- int iterator_fail_after_ = -1;
- std::atomic<int> delete_calls_ {0};
- std::atomic<size_t> total_keys_deleted_ {0};
- int fail_delete_after_ = -1; // -1 means never fail
+ size_t batch_size_;
+ int iterator_fail_after_;
+ int delete_calls_ = 0;
+ bool fail_delete_ = false;
+ std::vector<std::string> deleted_keys_;
};
-class RecyclerBatchDeleteTest : public testing::Test {
-protected:
- void SetUp() override {
- thread_pool_ = std::make_shared<SimpleThreadPool>(4);
- thread_pool_->start();
- }
-
- void TearDown() override {
- if (thread_pool_) {
- thread_pool_->stop();
- }
- }
-
- std::vector<ObjectMeta> generate_objects(size_t count) {
- std::vector<ObjectMeta> objects;
- objects.reserve(count);
- for (size_t i = 0; i < count; ++i) {
- objects.push_back(ObjectMeta {
- .key = "test_key_" + std::to_string(i),
- .size = 100,
- .mtime_s = 0,
- });
+class CountingRateLimitPolicy final : public ObjStorageRateLimitPolicy {
+public:
+ CountingRateLimitPolicy(size_t* get_requests, size_t* put_requests)
+ : get_requests_(get_requests), put_requests_(put_requests) {}
+
+ ObjStorageRateLimitToken acquire(ObjStorageRequestType type, size_t) const
override {
+ if (type == ObjStorageRequestType::GET) {
+ ++*get_requests_;
+ } else {
+ ++*put_requests_;
}
- return objects;
+ return {};
}
- std::shared_ptr<SimpleThreadPool> thread_pool_;
+private:
+ size_t* get_requests_;
+ size_t* put_requests_;
};
-// Test 1: Basic batch processing with multiple batches
-TEST_F(RecyclerBatchDeleteTest, MultipleBatches) {
- // Save original config and set small batch size for testing
- int32_t original_config = config::recycler_max_tasks_per_batch;
- config::recycler_max_tasks_per_batch = 3; // 3 tasks per batch
-
- // Create 10 objects, with batch_size=2 (keys per task),
max_tasks_per_batch=3
- // Expected: 10 objects / 2 keys per task = 5 tasks
- // 5 tasks / 3 tasks per batch = 2 batches (3 tasks + 2 tasks)
- auto objects = generate_objects(10);
- MockObjStorageClient client(objects);
-
- ObjClientOptions options;
- options.executor = thread_pool_;
-
- // Use batch_size=2 to create more tasks
- auto response = client.delete_objects_recursively_(
- {.bucket = "test_bucket", .key = "test_prefix"}, options, 0, 2);
-
- EXPECT_EQ(response.ret, 0);
- EXPECT_EQ(client.get_delete_calls(), 5); // 10 objects / 2 = 5
delete calls
- EXPECT_EQ(client.get_total_keys_deleted(), 10); // All 10 keys deleted
-
- // Restore config
- config::recycler_max_tasks_per_batch = original_config;
-}
-
-// Test 2: Iterator becomes invalid during iteration
-TEST_F(RecyclerBatchDeleteTest, IteratorInvalidMidway) {
- int32_t original_config = config::recycler_max_tasks_per_batch;
- config::recycler_max_tasks_per_batch = 100;
-
- // Create 20 objects but iterator fails after 10
- auto objects = generate_objects(20);
- MockObjStorageClient client(objects, 10); // fail_after=10
-
- ObjClientOptions options;
- options.executor = thread_pool_;
-
- auto response = client.delete_objects_recursively_(
- {.bucket = "test_bucket", .key = "test_prefix"}, options, 0, 5);
-
- // Should return error because iterator became invalid
- EXPECT_EQ(response.ret, -1);
- // Should have processed some objects before failure
- EXPECT_GT(client.get_total_keys_deleted(), 0);
- EXPECT_LT(client.get_total_keys_deleted(), 20);
-
- config::recycler_max_tasks_per_batch = original_config;
-}
-
-// Test 3: Delete operation fails (triggers cancel)
-TEST_F(RecyclerBatchDeleteTest, DeleteFailureTriggersCancel) {
- int32_t original_config = config::recycler_max_tasks_per_batch;
- config::recycler_max_tasks_per_batch = 10;
-
- auto objects = generate_objects(30);
- MockObjStorageClient client(objects);
- client.set_fail_delete_after(2); // Fail after 2 successful deletes
-
- ObjClientOptions options;
- options.executor = thread_pool_;
-
- auto response = client.delete_objects_recursively_(
- {.bucket = "test_bucket", .key = "test_prefix"}, options, 0, 5);
-
- // Should return error because delete failed
- EXPECT_EQ(response.ret, -1);
-
- config::recycler_max_tasks_per_batch = original_config;
-}
-
-// Test 4: Empty object list
-TEST_F(RecyclerBatchDeleteTest, EmptyObjectList) {
- int32_t original_config = config::recycler_max_tasks_per_batch;
- config::recycler_max_tasks_per_batch = 100;
-
- std::vector<ObjectMeta> empty_objects;
- MockObjStorageClient client(empty_objects);
-
- ObjClientOptions options;
- options.executor = thread_pool_;
-
- auto response = client.delete_objects_recursively_(
- {.bucket = "test_bucket", .key = "test_prefix"}, options, 0, 1000);
-
- EXPECT_EQ(response.ret, 0);
- EXPECT_EQ(client.get_delete_calls(), 0);
- EXPECT_EQ(client.get_total_keys_deleted(), 0);
-
- config::recycler_max_tasks_per_batch = original_config;
-}
-
-// Test 5: Objects less than batch_size
-TEST_F(RecyclerBatchDeleteTest, ObjectsLessThanBatchSize) {
- int32_t original_config = config::recycler_max_tasks_per_batch;
- config::recycler_max_tasks_per_batch = 100;
-
- auto objects = generate_objects(5);
- MockObjStorageClient client(objects);
-
- ObjClientOptions options;
- options.executor = thread_pool_;
-
- // batch_size=1000, but only 5 objects
- auto response = client.delete_objects_recursively_(
- {.bucket = "test_bucket", .key = "test_prefix"}, options, 0, 1000);
-
- EXPECT_EQ(response.ret, 0);
- EXPECT_EQ(client.get_delete_calls(), 1); // All 5 keys in one delete call
- EXPECT_EQ(client.get_total_keys_deleted(), 5);
-
- config::recycler_max_tasks_per_batch = original_config;
-}
-
-// Test 6: Exact batch boundary
-TEST_F(RecyclerBatchDeleteTest, ExactBatchBoundary) {
- int32_t original_config = config::recycler_max_tasks_per_batch;
- config::recycler_max_tasks_per_batch = 2; // 2 tasks per batch
-
- // 8 objects with batch_size=2 = 4 tasks
- // 4 tasks with max_tasks_per_batch=2 = exactly 2 batches
- auto objects = generate_objects(8);
- MockObjStorageClient client(objects);
-
- ObjClientOptions options;
- options.executor = thread_pool_;
-
- auto response = client.delete_objects_recursively_(
- {.bucket = "test_bucket", .key = "test_prefix"}, options, 0, 2);
-
- EXPECT_EQ(response.ret, 0);
- EXPECT_EQ(client.get_delete_calls(), 4); // 8 / 2 = 4 tasks
- EXPECT_EQ(client.get_total_keys_deleted(), 8);
-
- config::recycler_max_tasks_per_batch = original_config;
-}
-
-// Test 7: Invalid config value (negative)
-TEST_F(RecyclerBatchDeleteTest, InvalidConfigNegative) {
- int32_t original_config = config::recycler_max_tasks_per_batch;
- config::recycler_max_tasks_per_batch = -1; // Invalid negative value
-
- auto objects = generate_objects(10);
- MockObjStorageClient client(objects);
-
- ObjClientOptions options;
- options.executor = thread_pool_;
-
- // Should use default value 1000 and still work
- auto response = client.delete_objects_recursively_(
- {.bucket = "test_bucket", .key = "test_prefix"}, options, 0, 5);
-
- EXPECT_EQ(response.ret, 0);
- EXPECT_EQ(client.get_total_keys_deleted(), 10);
-
- config::recycler_max_tasks_per_batch = original_config;
-}
-
-// Test 8: Invalid config value (zero)
-TEST_F(RecyclerBatchDeleteTest, InvalidConfigZero) {
- int32_t original_config = config::recycler_max_tasks_per_batch;
- config::recycler_max_tasks_per_batch = 0; // Invalid zero value
-
- auto objects = generate_objects(10);
- MockObjStorageClient client(objects);
-
- ObjClientOptions options;
- options.executor = thread_pool_;
-
- // Should use default value 1000 and still work
- auto response = client.delete_objects_recursively_(
- {.bucket = "test_bucket", .key = "test_prefix"}, options, 0, 5);
-
- EXPECT_EQ(response.ret, 0);
- EXPECT_EQ(client.get_total_keys_deleted(), 10);
-
- config::recycler_max_tasks_per_batch = original_config;
-}
-
-// Test 9: Expiration time filtering
-TEST_F(RecyclerBatchDeleteTest, ExpirationTimeFiltering) {
- int32_t original_config = config::recycler_max_tasks_per_batch;
- config::recycler_max_tasks_per_batch = 100;
-
+std::vector<ObjectMeta> make_objects(size_t count) {
std::vector<ObjectMeta> objects;
- // Create 10 objects: 5 with old mtime (should be deleted), 5 with new
mtime (should be kept)
- for (int i = 0; i < 5; ++i) {
- objects.push_back(ObjectMeta {
- .key = "old_key_" + std::to_string(i),
- .size = 100,
- .mtime_s = 100, // Old timestamp
- });
- }
- for (int i = 0; i < 5; ++i) {
- objects.push_back(ObjectMeta {
- .key = "new_key_" + std::to_string(i),
+ for (size_t i = 0; i < count; ++i) {
+ objects.push_back({
+ .file_path = "test_key_" + std::to_string(i),
.size = 100,
- .mtime_s = 1000, // New timestamp
+ .mtime_s = static_cast<int64_t>(i),
});
}
+ return objects;
+}
- MockObjStorageClient client(objects);
+TEST(RecyclerBatchDeleteTest, UsesProviderBatchCapability) {
+ auto backend = std::make_shared<MockObjStorageBackend>(make_objects(10),
3);
+ ObjStorageClient client(backend);
+ auto response = client.delete_objects_recursively({.bucket = "bucket",
.prefix = "test_key_"});
Review Comment:
[P2] Exercise the production recursive-delete executor
All replacement tests call delete_objects_recursively without
RecursiveDeleteOptions, so they only run the sequential fallback. Production
S3Accessor always supplies the SyncExecutor-based pool path, including
max_tasks_per_batch, peer cancellation, and unfinished-task propagation; the
deleted fixture covered that branch with a real pool. Please retain a test
wired like make_recursive_delete_options that forces multiple task batches and
a delete failure, then verifies cancellation/error propagation (including
zero/negative recycler_max_tasks_per_batch handling).
##########
be/test/io/fs/s3_obj_storage_client_test.cpp:
##########
@@ -87,10 +93,15 @@ TEST_F(S3ObjStorageClientTest, put_list_delete_object) {
EXPECT_EQ(response.status.code, ErrorCode::OK);
// clang-format off
- response =
S3ObjStorageClientTest::obj_storage_client->list_objects({.bucket = bucket,
- .prefix = "S3ObjStorageClientTest/put_list_delete_object",},
&files);
+ iter = ObjectListIterator(S3ObjStorageClientTest::obj_storage_client,
{.bucket = bucket,
+ .key = "S3ObjStorageClientTest/put_list_delete_object"});
// clang-format on
- EXPECT_EQ(response.status.code, ErrorCode::OK);
+ for (auto obj = iter.next(); obj.results_.has_value(); obj = iter.next()) {
Review Comment:
[P2] Assert the iterator's terminal status
ObjectListIterator::next() uses an empty results_ for both clean EOF and a
failed list page. Because this loop only checks resp while a result exists, a
first-request failure skips the body and makes the expected-empty post-delete
checks pass. The old tests asserted the list status separately; please inspect
the terminal response or assert iter.is_valid() after every migrated loop
(including the role and Azure variants) so provider/rate-limit errors cannot
look like successful empty listings.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]