From 7946c7cb53f59e3fc14d2b62ff197ed859845df0 Mon Sep 17 00:00:00 2001 From: Mikhail Filimonov Date: Mon, 5 Oct 2026 17:33:51 +0200 Subject: [PATCH] Cut CAS GC batch deletes to the storage's batch-delete limit GC sent a write-once cohort (up to `gc_bulk_delete_chunk_keys`, 1000) as one `removeManyWriteOnce`, ignoring `objects_chunk_size_to_delete` on S3 and the single-key reality of storages without `DeleteObjects`. Add `IObjectStorage::batchDeleteKeyLimit` and `Backend::bulkDeleteKeyLimit` (forwarded by every decorator). `Gc::bulkDeleteChunkKeys` takes the smaller of the setting and that limit; `removeCohortWriteOnce` cuts the `manifest_deletes` and `cleanupRefObjects` cohorts into requests of that size. The cohort, its budget accounting and the one `authorityHolds` read per cohort are unchanged. The S3 delete paths also pass the store's error name into `S3Exception`: `isRefreshableCredentialError` recognises SDK-unmodelled errors only by `getExceptionName`. Tests: `CASBulkDeleteBackend` limit forwarding and clamping, GC cut for `manifest_deletes` and ref cleanup (no extra `gc/state` reads), and `CASS3BatchDelete` for the S3 limit and error names. Co-Authored-By: Claude Fable 5.1 Signed-off-by: Mikhail Filimonov --- docs/en/antalya/cas/configuration.md | 2 +- .../ContentAddressed/Backend/CasBackend.h | 4 + .../Backend/CasInMemoryBackend.h | 4 + .../Backend/CasInstrumentedBackend.h | 1 + .../Backend/CasObjectStorageBackend.cpp | 6 + .../Backend/CasObjectStorageBackend.h | 4 + .../Backend/CasThrottlingBackend.h | 1 + .../ContentAddressedSettings.cpp | 2 +- .../ContentAddressed/Gc/CasGc.cpp | 30 ++++- .../ContentAddressed/Gc/CasGc.h | 4 + .../ObjectStorages/IObjectStorage.h | 5 + .../ObjectStorages/S3/S3ObjectStorage.cpp | 20 +++- .../ObjectStorages/S3/S3ObjectStorage.h | 3 + src/Disks/tests/cas_test_helpers.h | 71 +++++++++++ .../tests/gtest_cas_bulk_delete_backend.cpp | 46 +++++++ src/Disks/tests/gtest_cas_decommission.cpp | 1 + .../gtest_cas_gc_manifest_bulk_delete.cpp | 31 +++++ src/Disks/tests/gtest_cas_mount.cpp | 1 + src/Disks/tests/gtest_cas_part_write.cpp | 5 + src/Disks/tests/gtest_cas_pool.cpp | 4 + src/Disks/tests/gtest_cas_ref_gc.cpp | 88 ++++++++++++++ .../gtest_cas_s3_bulk_delete_fallback.cpp | 112 +++++++++++++++++- src/IO/S3/deleteFileFromS3.cpp | 6 +- 23 files changed, 440 insertions(+), 11 deletions(-) diff --git a/docs/en/antalya/cas/configuration.md b/docs/en/antalya/cas/configuration.md index 17b1e3482768..a490428b3032 100644 --- a/docs/en/antalya/cas/configuration.md +++ b/docs/en/antalya/cas/configuration.md @@ -179,7 +179,7 @@ for the remaining caps, `0` means unbounded. |---|---|---|---| | `cas_manifest_sweep_list_budget_keys` | `1000` | `UInt64` | Orphan-manifest sweep `LIST` budget per round | | `cas_manifest_sweep_delete_budget_keys` | `100` | `UInt64` | Orphan-manifest sweep `DELETE` budget per round | -| `cas_gc_bulk_delete_chunk_keys` | `1000` | `1`–`1000` | Keys per batch delete request in GC's write-once families (owner-removed manifest bodies, covered ref logs and snapshots) | +| `cas_gc_bulk_delete_chunk_keys` | `1000` | `1`–`1000` | Keys per batch delete request in GC's write-once families (owner-removed manifest bodies, covered ref logs and snapshots), capped by the storage's batch-delete limit | | `cas_gc_round_graduation_budget` | `5000` | `0` = unbounded | Blob-graduation (`condemned` → `delete_pending`) cohort cap per round | | `cas_gc_round_redelete_budget` | `5000` | `0` = unbounded | Exact-token re-delete cohort cap for prior `delete_pending` rows per round | | `cas_gc_round_sweep_namespace_budget` | `20` | `0` = unbounded | Distinct namespaces per orphan-manifest sweep page whose protection view may be built | diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasBackend.h b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasBackend.h index 96911defc15f..eb502f9724c2 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasBackend.h +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasBackend.h @@ -201,6 +201,10 @@ class Backend /// absence is success. virtual void removeManyWriteOnce(const std::vector & keys, TransportAccess &) = 0; + /// The most keys one `removeManyWriteOnce` sends to the storage as one request; at least 1. A decorator + /// must forward it: the default of 1 would turn every batch below it into one-key requests. + virtual size_t bulkDeleteKeyLimit() const { return 1; } + /// Authoritative, cache-bypassing probe of one key -- see `ProbeOutcome`. DEFAULT (used by every /// backend without sharper raw-error evidence, e.g. `InMemoryBackend`): derived from `head`/`read` /// alone, so it can only distinguish `Present` from `KeyAbsent`, and ANY exception from either diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasInMemoryBackend.h b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasInMemoryBackend.h index 2074750e3ef6..9e3100efdf3e 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasInMemoryBackend.h +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasInMemoryBackend.h @@ -1,5 +1,6 @@ #pragma once #include +#include #include #include #include @@ -59,6 +60,9 @@ class InMemoryBackend : public Backend /// How many `removeManyWriteOnce` calls reached the store, armed failures included. size_t bulkRemoveCalls() const; + /// The emulation deletes keys one by one under its lock, so the CAS maximum is the only bound. + size_t bulkDeleteKeyLimit() const override { return kBulkDeleteMaxKeys; } + /// Creates the key when `expected_value` is empty, or replaces the incarnation it names. A /// refused precondition leaves the store unchanged. Value enforcement can be disabled with /// `setEnforceTokens` to model a backend that incorrectly ignores the condition. diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasInstrumentedBackend.h b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasInstrumentedBackend.h index ed0c25dab861..5d1f3f36e953 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasInstrumentedBackend.h +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasInstrumentedBackend.h @@ -161,6 +161,7 @@ class InstrumentedBackend final : public Backend bool supportsListTokens() const override { return inner->supportsListTokens(); } uint64_t attemptTimeoutMs() const override { return inner->attemptTimeoutMs(); } uint64_t attemptEnvelopeMs() const override { return inner->attemptEnvelopeMs(); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } bool refreshCredentials() override { return inner->refreshCredentials(); } private: diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasObjectStorageBackend.cpp b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasObjectStorageBackend.cpp index 0589929d2912..2b63dc0ef9d7 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasObjectStorageBackend.cpp +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasObjectStorageBackend.cpp @@ -1,6 +1,7 @@ #include #include +#include #include #include #include @@ -1003,6 +1004,11 @@ void ObjectStorageBackend::emuForgetDeletedToken(const String & key) } } +size_t ObjectStorageBackend::bulkDeleteKeyLimit() const +{ + return mode == Mode::Native ? object_storage->batchDeleteKeyLimit() : kBulkDeleteMaxKeys; +} + void ObjectStorageBackend::removeManyWriteOnce(const std::vector & keys, TransportAccess & access) { if (keys.empty()) diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasObjectStorageBackend.h b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasObjectStorageBackend.h index f9f51396efbb..cc3c83f4f143 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasObjectStorageBackend.h +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasObjectStorageBackend.h @@ -101,6 +101,10 @@ class ObjectStorageBackend final : public Backend /// under the control-plane profile; `EmulatedSingleProcess` deletes each present key under the /// emulation lock with the same token bookkeeping as the single-key delete. void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override; + + /// Native: the object storage's own limit. Emulated: `kBulkDeleteMaxKeys`, since this backend deletes + /// the keys one by one under its emulation lock. + size_t bulkDeleteKeyLimit() const override; /// Native mints its store's own dialect (ETag or GCS generation); the emulated adapter mints its /// own values. Dialect dialect() const override { return mode == Mode::Native ? native_token_type : Dialect::Emulated; } diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasThrottlingBackend.h b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasThrottlingBackend.h index e620195a31c3..8f69d4464347 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasThrottlingBackend.h +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasThrottlingBackend.h @@ -125,6 +125,7 @@ class ThrottlingBackend final : public Backend bool supportsListTokens() const override { return inner->supportsListTokens(); } uint64_t attemptTimeoutMs() const override { return inner->attemptTimeoutMs(); } uint64_t attemptEnvelopeMs() const override { return inner->attemptEnvelopeMs(); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } bool refreshCredentials() override { return inner->refreshCredentials(); } void checkPoolPreconditions() override { inner->checkPoolPreconditions(); } void checkSkipAccessCheckSupport() override { inner->checkSkipAccessCheckSupport(); } diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedSettings.cpp b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedSettings.cpp index 7331d2bad29a..63e3fbd70b14 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedSettings.cpp +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/ContentAddressedSettings.cpp @@ -80,7 +80,7 @@ constexpr std::string_view CAS_KEY_PREFIX = "cas_"; DECLARE(UInt64, manifest_decode_cache_bytes, 128ULL << 20, "Manifest DECODE cache byte budget (0 disables)", 0) \ DECLARE(UInt64, gc_meta_pool_size, 16, "Bounded pool size for GC per-hash freshness-meta writes", 0) \ DECLARE(UInt64, gc_io_concurrency, 16, "Maximum number of threads in the GC I/O pool. Used for fold and rebuild read-ahead, orphan-manifest sweep planning reads, and pending_deletes HEAD plus conditional DELETE. Per-hash meta writes use gc_meta_pool_size; other GC requests run on the round thread. 1 disables parallel GC I/O", 0) \ - DECLARE(UInt64, gc_bulk_delete_chunk_keys, 1000, "Keys per batch delete request in GC's write-once families (owner-removed manifest bodies, covered ref logs and snapshots); 1 to 1000", 0) \ + DECLARE(UInt64, gc_bulk_delete_chunk_keys, 1000, "Keys per batch delete request in GC's write-once families (owner-removed manifest bodies, covered ref logs and snapshots), capped by the storage's batch-delete limit; 1 to 1000", 0) \ DECLARE(UInt64, attempt_timeout_ms, 5000, "Budget for one HTTP attempt of a writable Native mount's control-plane requests (read, head, list, remove, conditional write), at least 1. With the connect cap it forms the attempt envelope the lease arithmetic reserves", 0) \ DECLARE(UInt64, lease_safety_margin_ms, 2000, "Startup-only margin validated against the mount lease TTL: attempt envelope + this must be strictly less than the TTL, and renew period + 2 × envelope + this too", 0) \ DECLARE(String, staging_backend, "local", "Blob staging backend (local | s3); s3 is opt-in", 0) \ diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.cpp b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.cpp index f871182c8edb..de10d48954a6 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.cpp +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.cpp @@ -408,6 +408,30 @@ uint64_t removeChunkWriteOnceOrOneByOne(CasOperation & op, const std::vector & cohort, size_t request_keys, const Retry & policy) +{ + uint64_t calls = 0; + for (size_t begin = 0; begin < cohort.size(); begin += request_keys) + { + const size_t end = std::min(cohort.size(), begin + request_keys); + const std::vector request(cohort.begin() + begin, cohort.begin() + end); + calls += removeChunkWriteOnceOrOneByOne(op, request, policy); + } + return calls; +} +} + +size_t Gc::bulkDeleteChunkKeys() const +{ + const uint64_t configured = store->poolConfig().gc_bulk_delete_chunk_keys; + const uint64_t storage_limit = store->poolBackendPtr()->bulkDeleteKeyLimit(); + return std::clamp(std::min(configured, storage_limit), 1, kBulkDeleteMaxKeys); +} + Gc::RedeleteIo Gc::performRedeleteIo(const RetiredEntry & entry, const Layout & layout, CasOperation & op) { RedeleteIo io; @@ -1279,6 +1303,7 @@ RoundReport Gc::runRegularRound(std::function on_lease_acquired, bool al /// never per key, so the next round's fold sees that key as already gone. The etag the fold /// observed rides the event as information only. const size_t chunk_keys = std::clamp(store->poolConfig().gc_bulk_delete_chunk_keys, 1, kBulkDeleteMaxKeys); + const size_t request_keys = bulkDeleteChunkKeys(); uint64_t attempted = 0; uint64_t requests = 0; std::vector chunk; @@ -1290,7 +1315,7 @@ RoundReport Gc::runRegularRound(std::function on_lease_acquired, bool al /// A backend without `DeleteObjects` (GCS) falls back to one admitted delete per key here; /// `chunk_entries`' per-key bookkeeping below is unaffected either way -- it counts objects /// that are gone after this call returns, not how many requests it took to get them there. - requests += removeChunkWriteOnceOrOneByOne(op, chunk, Retry::standard()); + requests += removeCohortWriteOnce(op, chunk, request_keys, Retry::standard()); for (const auto * entry : chunk_entries) { ++report.manifests_deleted; @@ -3834,6 +3859,7 @@ void Gc::cleanupRefObjects( } const size_t chunk_keys = std::clamp(store->poolConfig().gc_bulk_delete_chunk_keys, 1, kBulkDeleteMaxKeys); + const size_t request_keys = bulkDeleteChunkKeys(); for (size_t begin = 0; begin < cohort.size(); ) { /// Cumulative per-round cap in KEYS, exactly as before; a chunk is cut to what remains. The @@ -3850,7 +3876,7 @@ void Gc::cleanupRefObjects( /// A backend without `DeleteObjects` (GCS) falls back to one admitted delete per key here. /// The budget and the profile event below count OBJECTS in `chunk`, which is the same /// `chunk.size()` whichever way `removeChunkWriteOnceOrOneByOne` actually sent them. - removeChunkWriteOnceOrOneByOne(op, chunk, Retry::standard()); + removeCohortWriteOnce(op, chunk, request_keys, Retry::standard()); work_budget.ref_cleanup_objects_used += chunk.size(); ProfileEvents::increment(ProfileEvents::CASRefCleanupObjectsDeleted, chunk.size()); /// cleanup object deletion /// Advance by what was actually sent, not the nominal chunk size: the budget cap above can diff --git a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.h b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.h index 33f2bdcf09a8..48ddf5ef0786 100644 --- a/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.h +++ b/src/Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Gc/CasGc.h @@ -505,6 +505,10 @@ class Gc /// instance, and the scheduler holds `gc_round_mutex` across both the install and the round. void setPhaseSink(GcPhaseSink sink) { phase_sink = std::move(sink); } + /// Keys per GC write-once delete request: `gc_bulk_delete_chunk_keys` capped by the storage's batch + /// limit, clamped to [1, `kBulkDeleteMaxKeys`] because an empty request would never advance. + size_t bulkDeleteChunkKeys() const; + void setRebuildEdgeBudgetForTest(uint64_t n) { rebuild_edge_budget_override = n; } /// TEST SEAM: disable the round's journal trim so a folded event stays in the journal diff --git a/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.h b/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.h index 436f8c05bd51..00a2a0d21c44 100644 --- a/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.h +++ b/src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.h @@ -380,6 +380,11 @@ class IObjectStorage virtual void removeObjectsIfExistUnderProfile( const StoredObjects & objects, const ObjectStorageControlRequest & request); + /// The most objects one `removeObjectsIfExistUnderProfile` call sends as one request; at least 1. + /// Callers cut their batches to this size, so one call stays one storage request. The default + /// suits a storage without a batch delete. + virtual size_t batchDeleteKeyLimit() const { return 1; } + /// Copy object with different attributes if required virtual void copyObject( /// NOLINT const StoredObject & object_from, diff --git a/src/Disks/DiskObjectStorage/ObjectStorages/S3/S3ObjectStorage.cpp b/src/Disks/DiskObjectStorage/ObjectStorages/S3/S3ObjectStorage.cpp index 6eb155dd4c44..da1343df6de0 100644 --- a/src/Disks/DiskObjectStorage/ObjectStorages/S3/S3ObjectStorage.cpp +++ b/src/Disks/DiskObjectStorage/ObjectStorages/S3/S3ObjectStorage.cpp @@ -748,12 +748,15 @@ void S3ObjectStorage::removeObjectsIfExistImpl( throw Exception(ErrorCodes::NOT_IMPLEMENTED, "{} does not support DeleteObjects", getName()); } - throw S3Exception(err.GetErrorType(), "{} (Code: {}) while removing {} objects from S3 in one request", - err.GetMessage(), static_cast(err.GetErrorType()), objects.size()); + throw S3Exception( + PreformattedMessage::create("{} (Code: {}) while removing {} objects from S3 in one request", + err.GetMessage(), static_cast(err.GetErrorType()), objects.size()), + err.GetErrorType(), err.GetExceptionName()); } String failed_keys; std::optional first_error_type; + String first_error_name; for (const auto & err : outcome.GetResult().GetErrors()) { const auto error_type = classifyDeleteObjectsErrorCode(err.GetCode()); @@ -763,10 +766,21 @@ void S3ObjectStorage::removeObjectsIfExistImpl( failed_keys += ", "; failed_keys += err.GetKey() + " (" + err.GetCode() + ": " + err.GetMessage() + ")"; if (!first_error_type) + { first_error_type = error_type; + first_error_name = err.GetCode(); + } } if (first_error_type) - throw S3Exception(*first_error_type, "batch removal left objects behind: [{}]", failed_keys); + throw S3Exception( + PreformattedMessage::create("batch removal left objects behind: [{}]", failed_keys), *first_error_type, first_error_name); +} + +size_t S3ObjectStorage::batchDeleteKeyLimit() const +{ + if (const auto supported = s3_capabilities.isBatchDeleteSupported(); supported.has_value() && !*supported) + return 1; + return std::max(1, s3_settings.get()->request_settings[S3RequestSetting::objects_chunk_size_to_delete]); } bool S3ObjectStorage::conditionalOpsUseGenerationTokens() const diff --git a/src/Disks/DiskObjectStorage/ObjectStorages/S3/S3ObjectStorage.h b/src/Disks/DiskObjectStorage/ObjectStorages/S3/S3ObjectStorage.h index 5f5903f2c6aa..f11984dce135 100644 --- a/src/Disks/DiskObjectStorage/ObjectStorages/S3/S3ObjectStorage.h +++ b/src/Disks/DiskObjectStorage/ObjectStorages/S3/S3ObjectStorage.h @@ -137,6 +137,9 @@ class S3ObjectStorage : public IObjectStorage void removeObjectsIfExistUnderProfile( const StoredObjects & objects, const ObjectStorageControlRequest & request) override; + /// `objects_chunk_size_to_delete`, at least 1, while `DeleteObjects` is not known unsupported; 1 after. + size_t batchDeleteKeyLimit() const override; + void tagObjects(const StoredObjects & objects, const std::string & tag_key, const std::string & tag_value) override; ObjectMetadata getObjectMetadata(const std::string & path, bool with_tags) const override; diff --git a/src/Disks/tests/cas_test_helpers.h b/src/Disks/tests/cas_test_helpers.h index ae87f129fce3..5b02caedc9df 100644 --- a/src/Disks/tests/cas_test_helpers.h +++ b/src/Disks/tests/cas_test_helpers.h @@ -74,6 +74,7 @@ namespace DB::ContentAddressedSetting namespace DB::ErrorCodes { extern const int CORRUPTED_DATA; + extern const int NOT_IMPLEMENTED; } namespace DB::Cas::tests @@ -1972,6 +1973,76 @@ class CountingBackend : public DB::Cas::InMemoryBackend uint64_t publish_total = 0; }; +/// The in-memory store with `S3ObjectStorage`'s batch-delete capability. A `removeManyWriteOnce` of more +/// than one key is refused locally with `NOT_IMPLEMENTED` once the capability is false; a store that rejects +/// `DeleteObjects` turns an unknown capability false on its first multi-key request. One key always goes. +class BatchCapabilityBackend : public CountingBackend +{ +public: + size_t bulkDeleteKeyLimit() const override + { + std::lock_guard lock(capability_mutex); + return batch_supported == std::optional{false} ? 1 : storage_limit; + } + + void removeManyWriteOnce(const std::vector & keys, DB::Cas::TransportAccess & access) override + { + { + std::lock_guard lock(capability_mutex); + call_sizes.push_back(keys.size()); + if (keys.size() > 1 && batch_supported == std::optional{false}) + throw DB::Exception(DB::ErrorCodes::NOT_IMPLEMENTED, "batch delete is known unsupported"); + request_sizes.push_back(keys.size()); + if (keys.size() > 1 && !batch_supported && rejects_batches) + { + batch_supported = false; + throw DB::Exception(DB::ErrorCodes::NOT_IMPLEMENTED, "the store rejected DeleteObjects"); + } + } + CountingBackend::removeManyWriteOnce(keys, access); + } + + void setBatchDeleteSupported(std::optional value) + { + std::lock_guard lock(capability_mutex); + batch_supported = value; + } + + void setStoreRejectsBatches(bool value) + { + std::lock_guard lock(capability_mutex); + rejects_batches = value; + } + + void setStorageLimit(size_t value) + { + std::lock_guard lock(capability_mutex); + storage_limit = value; + } + + /// Every `removeManyWriteOnce` that reached this backend, locally refused ones included. + std::vector callSizes() const + { + std::lock_guard lock(capability_mutex); + return call_sizes; + } + + /// Requests the modeled store received: calls minus local refusals. + std::vector requestSizes() const + { + std::lock_guard lock(capability_mutex); + return request_sizes; + } + +private: + mutable std::mutex capability_mutex; + std::optional batch_supported; + bool rejects_batches = false; + size_t storage_limit = DB::Cas::kBulkDeleteMaxKeys; + std::vector call_sizes; + std::vector request_sizes; +}; + /// Records the ORDER of writes (so a test can compare indices) and lets a test refuse or fail chosen /// writes by key. Delegates every request to `CountingBackend` unchanged, so the per-key counters /// remain available as the positive control. diff --git a/src/Disks/tests/gtest_cas_bulk_delete_backend.cpp b/src/Disks/tests/gtest_cas_bulk_delete_backend.cpp index 68aeabeb9523..b41851c8aa30 100644 --- a/src/Disks/tests/gtest_cas_bulk_delete_backend.cpp +++ b/src/Disks/tests/gtest_cas_bulk_delete_backend.cpp @@ -8,6 +8,8 @@ #include #include #include +#include +#include #include #include #include @@ -140,6 +142,14 @@ TEST(CASBulkDeleteBackend, InstrumentedCountsOneRequestAndOneDeletePerKeyClass) } #if USE_AWS_S3 +TEST(CASBulkDeleteBackend, ThrottlingBackendForwardsStorageLimit) +{ + auto inner = std::make_shared(); + inner->setStorageLimit(250); + ThrottlingBackend backend(inner, ThrottlingBackend::Mode::FirstPerKey, 1, 429); + EXPECT_EQ(backend.bulkDeleteKeyLimit(), 250u); +} + TEST(CASBulkDeleteBackend, ThrottlingRefusesTheChunkOnceAndTheEngineReissuesIt) { auto inner = std::make_shared(); @@ -195,3 +205,39 @@ TEST(CASBulkDeleteBackend, LocalObjectStorageRefusesTheProfileOverload) }); } #endif + +TEST(CASBulkDeleteBackend, InMemoryStorageLimitDefaultsToTheCasMaximum) +{ + InMemoryBackend backend; + EXPECT_EQ(backend.bulkDeleteKeyLimit(), kBulkDeleteMaxKeys); +} + +TEST(CASBulkDeleteBackend, PoolWrappedBackendForwardsStorageLimit) +{ + auto backend = std::make_shared(); + backend->setStorageLimit(250); + auto store = DB::Cas::tests::openPoolForTest(backend, /*gc_fold_max_defer_rounds=*/ 0); + ASSERT_NE(dynamic_cast(store->poolBackendPtr().get()), nullptr) + << "the pool must wrap its backend, or this test proves nothing about the decorator"; + EXPECT_EQ(store->poolBackendPtr()->bulkDeleteKeyLimit(), 250u); + Gc gc(store, DB::UInt128{1}); + EXPECT_EQ(gc.bulkDeleteChunkKeys(), 250u); +} + +TEST(CASBulkDeleteBackend, GcChunkIsTheSmallerOfSettingAndStorageLimit) +{ + auto backend = std::make_shared(); + backend->setStorageLimit(700); + auto store = Pool::open(backend, PoolConfig{.pool_prefix = "p", .server_root_id = "test", .gc_bulk_delete_chunk_keys = 300}); + EXPECT_EQ(Gc(store, DB::UInt128{1}).bulkDeleteChunkKeys(), 300u); + backend->setBatchDeleteSupported(false); + EXPECT_EQ(Gc(store, DB::UInt128{1}).bulkDeleteChunkKeys(), 1u); +} + +TEST(CASBulkDeleteBackend, ZeroStorageLimitClampsGcChunkToOneKey) +{ + auto backend = std::make_shared(); + backend->setStorageLimit(0); + auto store = DB::Cas::tests::openPoolForTest(backend, 0); + EXPECT_EQ(Gc(store, DB::UInt128{1}).bulkDeleteChunkKeys(), 1u); +} diff --git a/src/Disks/tests/gtest_cas_decommission.cpp b/src/Disks/tests/gtest_cas_decommission.cpp index d590eaed7a59..936e9655e9b1 100644 --- a/src/Disks/tests/gtest_cas_decommission.cpp +++ b/src/Disks/tests/gtest_cas_decommission.cpp @@ -1329,6 +1329,7 @@ class FailDeletesUnderPrefixBackend : public Backend return inner->remove(key, expected_value, access); } void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override { inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { diff --git a/src/Disks/tests/gtest_cas_gc_manifest_bulk_delete.cpp b/src/Disks/tests/gtest_cas_gc_manifest_bulk_delete.cpp index 3a7c2ca3ba7f..edbf77504839 100644 --- a/src/Disks/tests/gtest_cas_gc_manifest_bulk_delete.cpp +++ b/src/Disks/tests/gtest_cas_gc_manifest_bulk_delete.cpp @@ -7,6 +7,8 @@ #include #include +#include + /// The manifest_deletes phase sends owner-removed manifest bodies to the store in chunks of /// write-once keys, one request per chunk, and records every chunk that succeeded before a later /// one can fail. @@ -194,3 +196,32 @@ TEST(CASGCManifestBulkDelete, ASuppressedRoundMakesNoRequest) for (const ManifestId & id : ids) EXPECT_TRUE((*op).head(store->layout().manifestKey(id), Retry::once()).has_value()); } + +/// The storage's batch limit caps every request below the configured chunk size: 8 bodies under a +/// limit of 3 go out as 3 + 3 + 2. +TEST(CASGCManifestBulkDelete, StorageLimitCutsTheCohortIntoRequestsOfAtMostThatManyKeys) +{ + auto backend = std::make_shared(); + backend->setStorageLimit(3); + auto store = Pool::open(backend, PoolConfig{.pool_prefix = "p", .server_root_id = "test", .gc_fold_max_defer_rounds = 0}); + const auto ids = seedDroppedManifests(*backend, store->layout(), 8); + + uint64_t requests_metric = 0; + Gc gc(store, kGc); + gc.setPhaseSink([&](const GcPhaseRecord & rec) + { + if (rec.phase == "manifest_deletes") + requests_metric += rec.metrics.at("requests"); + }); + const uint64_t deleted = reclaim(gc, store, *backend, ids, 16); + + EXPECT_EQ(deleted, 8u); + EXPECT_EQ(requests_metric, 3u); + const std::vector sizes = backend->callSizes(); + ASSERT_FALSE(sizes.empty()); + EXPECT_LE(*std::max_element(sizes.begin(), sizes.end()), 3u); + EXPECT_EQ(sizes, (std::vector{3, 3, 2})); + OperationForTest op(*backend); + for (const ManifestId & id : ids) + EXPECT_FALSE((*op).head(store->layout().manifestKey(id), Retry::once()).has_value()); +} diff --git a/src/Disks/tests/gtest_cas_mount.cpp b/src/Disks/tests/gtest_cas_mount.cpp index a9472994237f..783bd94b1dae 100644 --- a/src/Disks/tests/gtest_cas_mount.cpp +++ b/src/Disks/tests/gtest_cas_mount.cpp @@ -1013,6 +1013,7 @@ class AlwaysVanishesBackend final : public DB::Cas::Backend RawListPage list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access) override { return inner->list(prefix, cursor, limit, access); } RawRemoval remove(const String & key, const String & expected_value, TransportAccess & access) override { return inner->remove(key, expected_value, access); } void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override { inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { diff --git a/src/Disks/tests/gtest_cas_part_write.cpp b/src/Disks/tests/gtest_cas_part_write.cpp index 7259f13360ca..461718040aac 100644 --- a/src/Disks/tests/gtest_cas_part_write.cpp +++ b/src/Disks/tests/gtest_cas_part_write.cpp @@ -204,6 +204,7 @@ class HeadThenDeleteOnceBackend final : public DB::Cas::Backend RawListPage list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access) override { return inner->list(prefix, cursor, limit, access); } RawRemoval remove(const String & key, const String & expected_value, TransportAccess & access) override { return inner->remove(key, expected_value, access); } void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override { inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { @@ -249,6 +250,7 @@ class KeyCountingBackend final : public DB::Cas::Backend RawListPage list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access) override { return inner->list(prefix, cursor, limit, access); } RawRemoval remove(const String & key, const String & expected_value, TransportAccess & access) override { return inner->remove(key, expected_value, access); } void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override { inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { @@ -979,6 +981,7 @@ TEST(CASPartWriteTxn, PutBlobCondemnedDedupNeverGetsTheDyingObject) RawListPage list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access) override { return inner->list(prefix, cursor, limit, access); } RawRemoval remove(const String & key, const String & expected_value, TransportAccess & access) override { return inner->remove(key, expected_value, access); } void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override { inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { @@ -1072,6 +1075,7 @@ TEST(CASPartWriteTxn, PutBlobCondemnedDedupPresentNeverGetsTheDyingObject) RawListPage list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access) override { return inner->list(prefix, cursor, limit, access); } RawRemoval remove(const String & key, const String & expected_value, TransportAccess & access) override { return inner->remove(key, expected_value, access); } void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override { inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { @@ -1670,6 +1674,7 @@ TEST(CASPartWriteTxn, AdoptEvidenceRecordsTrustedManifestDependencyProofWithoutI RawListPage list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access) override { return inner->list(prefix, cursor, limit, access); } RawRemoval remove(const String & key, const String & expected_value, TransportAccess & access) override { return inner->remove(key, expected_value, access); } void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override { inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { diff --git a/src/Disks/tests/gtest_cas_pool.cpp b/src/Disks/tests/gtest_cas_pool.cpp index 452511f5f761..9825cae00a86 100644 --- a/src/Disks/tests/gtest_cas_pool.cpp +++ b/src/Disks/tests/gtest_cas_pool.cpp @@ -83,6 +83,7 @@ class WriteCountingBackend final : public DB::Cas::Backend ++writes; inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { @@ -239,6 +240,7 @@ class ProbeWatchingBackend final : public DB::Cas::Backend note(key.str()); inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { @@ -315,6 +317,7 @@ class ForwardingBackend : public DB::Cas::Backend RawListPage list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access) override { return inner->list(prefix, cursor, limit, access); } RawRemoval remove(const String & key, const String & expected_value, TransportAccess & access) override { return inner->remove(key, expected_value, access); } void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override { inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { @@ -1502,6 +1505,7 @@ class FenceInAdoptWindowBackend final : public DB::Cas::Backend RawListPage list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access) override { return inner->list(prefix, cursor, limit, access); } RawRemoval remove(const String & key, const String & expected_value, TransportAccess & access) override { return inner->remove(key, expected_value, access); } void removeManyWriteOnce(const std::vector & keys, TransportAccess & access) override { inner->removeManyWriteOnce(keys, access); } + size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); } std::expected write(const String & key, const String & bytes, const std::optional & expected_value, TransportAccess & access) override { diff --git a/src/Disks/tests/gtest_cas_ref_gc.cpp b/src/Disks/tests/gtest_cas_ref_gc.cpp index 835ddb197cbe..903c9371eac4 100644 --- a/src/Disks/tests/gtest_cas_ref_gc.cpp +++ b/src/Disks/tests/gtest_cas_ref_gc.cpp @@ -11,6 +11,7 @@ #include +#include #include #include @@ -1425,3 +1426,90 @@ TEST(CASRefGc, CatalogAdmittedFreshLifeWithoutParentSeedsSuccessorSeal) op.read(layout.foldSealKey(state.snap_generation, state.snap_attempt), Retry::once())->bytes); EXPECT_TRUE(seal.ref_lives.contains(life_id)); } + +namespace +{ + +struct RefCleanupRun +{ + uint64_t cohort_size = 0; + uint64_t deleted_metric = 0; + uint64_t gc_state_reads = 0; + std::vector call_sizes; + bool every_planned_key_gone = true; +}; + +/// One regular round over `table_count` covered logs under a backend whose batch limit is `storage_limit`. +RefCleanupRun runRefCleanupUnderStorageLimit(size_t storage_limit, uint64_t table_count) +{ + auto backend = std::make_shared(); + backend->setStorageLimit(storage_limit); + auto store = openPoolForTest(backend, /*gc_fold_max_defer_rounds*/ 0); + const Layout & layout = store->layout(); + const RootNamespace ns{"00/aa@cas@"}; + fixture::admitLive(*backend, layout, ns); + + RefTableListing listing; + uint64_t last = 0; + for (uint64_t i = 1; i <= table_count; ++i) + { + const ManifestRef r = mref(i); + writeManifestRaw(*backend, layout, ns, r, {blobEntryFor("a", DB::UInt128(i))}); + last = publishCommittedTransition(*backend, layout, ns, "t" + std::to_string(i), std::nullopt, r); + listing.logs.push_back(RefTxnId{1, last}); + } + std::vector committed; + for (uint64_t i = 1; i <= table_count; ++i) + committed.push_back(committedRow("t" + std::to_string(i), mref(i))); + writeRefSnapshotRaw(*backend, layout, minimalLiveSnapshot(ns.string(), RefTxnId{1, last}, committed)); + listing.snapshots.push_back(RefTxnId{1, last}); + replaceRecoverableCkptForRawFixture(*backend, layout, ns, RefCkpt{ + .life_epoch = 1, + .committed_through = RefTxnId{1, last}, + .checkpoint_snapshot_id = RefTxnId{1, last}, + .last_epoch_seal = std::nullopt, + }); + + const RefCleanupPlan plan = planRefCleanup(listing, RefTxnId{1, last}, RefTxnId{1, last}, std::nullopt); + RefCleanupRun run; + run.cohort_size = plan.deletable_logs.size() + plan.deletable_snapshots.size(); + + const auto deleted_before = ProfileEvents::global_counters[ProfileEvents::CASRefCleanupObjectsDeleted].load(); + const uint64_t state_reads_before = backend->getCount(layout.gcStateKey()); + const size_t calls_before = backend->callSizes().size(); + Gc gc(store, kGc); + EXPECT_TRUE(runRegularRoundReclaiming(gc).acquired_lease); + run.deleted_metric = ProfileEvents::global_counters[ProfileEvents::CASRefCleanupObjectsDeleted].load() - deleted_before; + run.gc_state_reads = backend->getCount(layout.gcStateKey()) - state_reads_before; + const std::vector sizes = backend->callSizes(); + run.call_sizes.assign(sizes.begin() + calls_before, sizes.end()); + + const NamespaceLifeId life = fixture::fixtureLife(ns); + OperationForTest op(*backend); + for (const RefTxnId & id : plan.deletable_logs) + run.every_planned_key_gone &= !(*op).head(layout.refLogKey(life, id), Retry::once()).has_value(); + for (const RefTxnId & id : plan.deletable_snapshots) + run.every_planned_key_gone &= !(*op).head(layout.refSnapshotKey(life, id), Retry::once()).has_value(); + return run; +} + +} + +/// The storage's batch limit cuts the cleanup cohort into smaller requests, and the extra requests cost +/// no extra authority reads: `authorityHolds` runs once per cohort, not once per request. +TEST(CASRefGc, RefObjectCleanupCutsRequestsToTheStorageLimitWithoutExtraAuthorityReads) +{ + const RefCleanupRun limited = runRefCleanupUnderStorageLimit(3, 8); + const RefCleanupRun unlimited = runRefCleanupUnderStorageLimit(kBulkDeleteMaxKeys, 8); + + ASSERT_GT(limited.cohort_size, 3u) << "the cohort must exceed the limit for the cut to be observable"; + EXPECT_EQ(limited.cohort_size, unlimited.cohort_size); + EXPECT_TRUE(limited.every_planned_key_gone); + EXPECT_EQ(limited.deleted_metric, limited.cohort_size); + EXPECT_EQ(unlimited.deleted_metric, unlimited.cohort_size); + + ASSERT_FALSE(limited.call_sizes.empty()); + EXPECT_LE(*std::max_element(limited.call_sizes.begin(), limited.call_sizes.end()), 3u); + EXPECT_GT(limited.call_sizes.size(), unlimited.call_sizes.size()); + EXPECT_EQ(limited.gc_state_reads, unlimited.gc_state_reads); +} diff --git a/src/Disks/tests/gtest_cas_s3_bulk_delete_fallback.cpp b/src/Disks/tests/gtest_cas_s3_bulk_delete_fallback.cpp index c0e4ee5f5306..d5bbf95a8dba 100644 --- a/src/Disks/tests/gtest_cas_s3_bulk_delete_fallback.cpp +++ b/src/Disks/tests/gtest_cas_s3_bulk_delete_fallback.cpp @@ -22,6 +22,7 @@ #include #include #include +#include #include #include @@ -38,6 +39,11 @@ namespace DB::ErrorCodes extern const int NOT_IMPLEMENTED; } +namespace DB::S3RequestSetting +{ +extern const S3RequestSettingsUInt64 objects_chunk_size_to_delete; +} + /// `S3ObjectStorage::removeObjectsIfExistImpl` (the CAS bulk-delete path, reached through /// `removeObjectsIfExistUnderProfile`) must honour `S3Capabilities::isBatchDeleteSupported()` the same /// way the generic `deleteFilesFromS3` does, but WITHOUT looping over the objects itself: once the @@ -183,7 +189,18 @@ void sendSingleDeleteError(Poco::Net::HTTPServerResponse & response, Poco::Net:: sendXml(response, status, body); } -std::shared_ptr makeStorageForTest(const std::string & endpoint, const DB::S3Capabilities & capabilities) +/// A quiet-mode `DeleteObjects` reply (HTTP 200) listing one `` per `{key, code}` pair. +void sendDeleteObjectsResult(Poco::Net::HTTPServerResponse & response, const std::vector> & errors) +{ + std::string body = ""; + for (const auto & [key, code] : errors) + body += "" + key + "" + code + "" + code + ""; + body += ""; + sendXml(response, Poco::Net::HTTPResponse::HTTP_OK, body); +} + +std::shared_ptr makeStorageForTest( + const std::string & endpoint, const DB::S3Capabilities & capabilities, std::optional objects_chunk_size_to_delete = {}) { DB::RemoteHostFilter remote_host_filter; DB::S3::PocoHTTPClientConfiguration cfg = DB::S3::ClientFactory::instance().createClientConfiguration( @@ -215,8 +232,11 @@ std::shared_ptr makeStorageForTest(const std::string & endp .is_s3express_bucket = false, }, "ACCESS_KEY_ID", "SECRET_ACCESS_KEY", "", {}, {}, DB::S3::CredentialsConfiguration{}); + auto settings = std::make_unique(); + if (objects_chunk_size_to_delete) + settings->request_settings[DB::S3RequestSetting::objects_chunk_size_to_delete] = *objects_chunk_size_to_delete; return std::make_shared( - std::move(client), std::make_unique(), + std::move(client), std::move(settings), DB::S3::URI(endpoint + "/test-bucket/"), capabilities, DB::ObjectStorageKeyGeneratorPtr{}, "disk"); } @@ -399,4 +419,92 @@ TEST(S3BulkDeleteFallback, ExplicitlyDisabledCapabilityThrowsNotImplementedWitho EXPECT_EQ(server.countMethod("DELETE"), 0u); } +TEST(CASS3BatchDelete, KeyLimitFollowsChunkSettingAndDropsToOneOnceBatchDeleteIsUnsupported) +{ + (void)contextForTest(); + + ScriptedS3Server server([](const Poco::Net::HTTPServerRequest &, Poco::Net::HTTPServerResponse & response) + { + sendBatchNotImplemented(response); + }); + + EXPECT_EQ(makeStorageForTest(server.getUrl(), DB::S3Capabilities{})->batchDeleteKeyLimit(), 1000u); + EXPECT_EQ(makeStorageForTest(server.getUrl(), DB::S3Capabilities{}, 500)->batchDeleteKeyLimit(), 500u); + EXPECT_EQ(makeStorageForTest(server.getUrl(), DB::S3Capabilities{}, 0)->batchDeleteKeyLimit(), 1u); + EXPECT_EQ(makeStorageForTest(server.getUrl(), DB::S3Capabilities{false}, 500)->batchDeleteKeyLimit(), 1u); + + auto learning = makeStorageForTest(server.getUrl(), DB::S3Capabilities{}, 500); + expectNotImplemented([&] + { + learning->removeObjectsIfExistUnderProfile( + {DB::StoredObject("key-a"), DB::StoredObject("key-b")}, DB::ObjectStorageControlRequest{}); + }); + EXPECT_EQ(learning->batchDeleteKeyLimit(), 1u); +} + +namespace +{ + +void expectRefusalNamed(const std::function & remove, const std::string & name) +{ + try + { + remove(); + FAIL() << "expected an S3Exception named " << name; + } + catch (const DB::S3Exception & e) + { + EXPECT_EQ(e.getExceptionName(), name) << e.message(); + } +} + +const DB::StoredObjects kTwoObjects{DB::StoredObject("key-a"), DB::StoredObject("key-b")}; + +} + +TEST(CASS3BatchDelete, RefusedDeleteObjectsKeepsTheStoresErrorName) +{ + (void)contextForTest(); + + for (const std::string code : {"EntityTooLarge", "MalformedXML"}) + { + SCOPED_TRACE("whole-request " + code); + ScriptedS3Server server([code](const Poco::Net::HTTPServerRequest &, Poco::Net::HTTPServerResponse & response) + { + sendSingleDeleteError(response, Poco::Net::HTTPResponse::HTTP_BAD_REQUEST, code, "refused"); + }); + auto storage = makeStorageForTest(server.getUrl(), DB::S3Capabilities{}); + expectRefusalNamed([&] { storage->removeObjectsIfExistUnderProfile(kTwoObjects, DB::ObjectStorageControlRequest{}); }, code); + } + { + SCOPED_TRACE("per-key EntityTooLarge in a 200 reply"); + ScriptedS3Server server([](const Poco::Net::HTTPServerRequest &, Poco::Net::HTTPServerResponse & response) + { + sendDeleteObjectsResult(response, {{"key-b", "EntityTooLarge"}}); + }); + auto storage = makeStorageForTest(server.getUrl(), DB::S3Capabilities{}); + expectRefusalNamed([&] { storage->removeObjectsIfExistUnderProfile(kTwoObjects, DB::ObjectStorageControlRequest{}); }, "EntityTooLarge"); + } + { + /// `NoSuchKey` is skipped; the first real error names the exception. + SCOPED_TRACE("per-key order: NoSuchKey, AccessDenied, EntityTooLarge"); + ScriptedS3Server server([](const Poco::Net::HTTPServerRequest &, Poco::Net::HTTPServerResponse & response) + { + sendDeleteObjectsResult(response, {{"key-a", "NoSuchKey"}, {"key-b", "AccessDenied"}, {"key-c", "EntityTooLarge"}}); + }); + auto storage = makeStorageForTest(server.getUrl(), DB::S3Capabilities{}); + expectRefusalNamed([&] { storage->removeObjectsIfExistUnderProfile(kTwoObjects, DB::ObjectStorageControlRequest{}); }, "AccessDenied"); + } + { + SCOPED_TRACE("one-key EntityTooLarge through deleteFileFromS3"); + ScriptedS3Server server([](const Poco::Net::HTTPServerRequest &, Poco::Net::HTTPServerResponse & response) + { + sendSingleDeleteError(response, Poco::Net::HTTPResponse::HTTP_BAD_REQUEST, "EntityTooLarge", "refused"); + }); + auto storage = makeStorageForTest(server.getUrl(), DB::S3Capabilities{}); + expectRefusalNamed([&] { storage->removeObjectsIfExistUnderProfile({DB::StoredObject("solo")}, DB::ObjectStorageControlRequest{}); }, "EntityTooLarge"); + EXPECT_EQ(server.countMethod("DELETE"), 1u); + } +} + #endif diff --git a/src/IO/S3/deleteFileFromS3.cpp b/src/IO/S3/deleteFileFromS3.cpp index 0111e00534a8..900f2c8cd4c3 100644 --- a/src/IO/S3/deleteFileFromS3.cpp +++ b/src/IO/S3/deleteFileFromS3.cpp @@ -74,8 +74,10 @@ void deleteFileFromS3( else { const auto & err = outcome.GetError(); - throw S3Exception(err.GetErrorType(), "{} (Code: {}) while removing object with path {} from S3", - err.GetMessage(), static_cast(err.GetErrorType()), key); + throw S3Exception( + PreformattedMessage::create("{} (Code: {}) while removing object with path {} from S3", + err.GetMessage(), static_cast(err.GetErrorType()), key), + err.GetErrorType(), err.GetExceptionName()); } }