Skip to content
Open
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
2 changes: 1 addition & 1 deletion docs/en/antalya/cas/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -201,6 +201,10 @@ class Backend
/// absence is success.
virtual void removeManyWriteOnce(const std::vector<WriteOnceKey> & 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
Expand Down
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
#pragma once
#include <Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasBackend.h>
#include <Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasRequests.h>
#include <exception>
#include <functional>
#include <map>
Expand Down Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
#include <Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasObjectStorageBackend.h>

#include <Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasEtag.h>
#include <Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Backend/CasRequests.h>
#include <Disks/DiskObjectStorage/MetadataStorages/ContentAddressed/Formats/CasFormat.h>
#include <Disks/DiskObjectStorage/ObjectStorages/Local/LocalObjectStorage.h>
#include <Disks/DiskObjectStorage/ObjectStorages/ObjectStorageIterator.h>
Expand Down Expand Up @@ -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<WriteOnceKey> & keys, TransportAccess & access)
{
if (keys.empty())
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<WriteOnceKey> & 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; }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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(); }
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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) \
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -408,6 +408,30 @@ uint64_t removeChunkWriteOnceOrOneByOne(CasOperation & op, const std::vector<Wri
}
}

namespace
{
/// Deletes `cohort` as consecutive requests of at most `request_keys` keys; returns the call count.
uint64_t removeCohortWriteOnce(
CasOperation & op, const std::vector<WriteOnceKey> & 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<WriteOnceKey> 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<size_t>(std::min(configured, storage_limit), 1, kBulkDeleteMaxKeys);
}

Gc::RedeleteIo Gc::performRedeleteIo(const RetiredEntry & entry, const Layout & layout, CasOperation & op)
{
RedeleteIo io;
Expand Down Expand Up @@ -1279,6 +1303,7 @@ RoundReport Gc::runRegularRound(std::function<void()> 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<size_t>(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<WriteOnceKey> chunk;
Expand All @@ -1290,7 +1315,7 @@ RoundReport Gc::runRegularRound(std::function<void()> 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;
Expand Down Expand Up @@ -3834,6 +3859,7 @@ void Gc::cleanupRefObjects(
}

const size_t chunk_keys = std::clamp<size_t>(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
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
5 changes: 5 additions & 0 deletions src/Disks/DiskObjectStorage/ObjectStorages/IObjectStorage.h
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
20 changes: 17 additions & 3 deletions src/Disks/DiskObjectStorage/ObjectStorages/S3/S3ObjectStorage.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -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<size_t>(err.GetErrorType()), objects.size());
throw S3Exception(
PreformattedMessage::create("{} (Code: {}) while removing {} objects from S3 in one request",
err.GetMessage(), static_cast<size_t>(err.GetErrorType()), objects.size()),
err.GetErrorType(), err.GetExceptionName());
}

String failed_keys;
std::optional<Aws::S3::S3Errors> first_error_type;
String first_error_name;
for (const auto & err : outcome.GetResult().GetErrors())
{
const auto error_type = classifyDeleteObjectsErrorCode(err.GetCode());
Expand All @@ -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<size_t>(1, s3_settings.get()->request_settings[S3RequestSetting::objects_chunk_size_to_delete]);
}

bool S3ObjectStorage::conditionalOpsUseGenerationTokens() const
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
71 changes: 71 additions & 0 deletions src/Disks/tests/cas_test_helpers.h
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,7 @@ namespace DB::ContentAddressedSetting
namespace DB::ErrorCodes
{
extern const int CORRUPTED_DATA;
extern const int NOT_IMPLEMENTED;
}

namespace DB::Cas::tests
Expand Down Expand Up @@ -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<bool>{false} ? 1 : storage_limit;
}

void removeManyWriteOnce(const std::vector<DB::Cas::WriteOnceKey> & 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<bool>{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<bool> 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<size_t> callSizes() const
{
std::lock_guard lock(capability_mutex);
return call_sizes;
}

/// Requests the modeled store received: calls minus local refusals.
std::vector<size_t> requestSizes() const
{
std::lock_guard lock(capability_mutex);
return request_sizes;
}

private:
mutable std::mutex capability_mutex;
std::optional<bool> batch_supported;
bool rejects_batches = false;
size_t storage_limit = DB::Cas::kBulkDeleteMaxKeys;
std::vector<size_t> call_sizes;
std::vector<size_t> 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.
Expand Down
Loading
Loading