Skip to content
Draft
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
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,11 @@
#include <variant>
#include <vector>

namespace DB::ErrorCodes
{
extern const int NOT_IMPLEMENTED;
}

namespace DB::Cas
{

Expand Down Expand Up @@ -169,6 +174,11 @@ class Backend
/// `DeleteMarker` is a removal that did NOT reclaim: a versioned bucket archived a noncurrent
/// version instead. Distinct from `Removed` because reclaiming the storage is the point.
enum class RawRemoval : uint8_t { Removed, Gone, Mismatch, DeleteMarker };
struct RawConditionalRemove
{
String key;
String expected_token;
};

/// Reads the whole object, or nullopt when the key is absent.
virtual std::optional<Raw> read (const String & key, TransportAccess &) = 0;
Expand Down Expand Up @@ -201,10 +211,23 @@ class Backend
/// absence is success.
virtual void removeManyWriteOnce(const std::vector<WriteOnceKey> & keys, TransportAccess &) = 0;

virtual std::vector<RawRemoval> removeManyByConditional(
const std::vector<RawConditionalRemove> & removals, TransportAccess & access)
{
if (removals.size() > 1)
throw Exception(ErrorCodes::NOT_IMPLEMENTED, "Conditional batch removal is not implemented by this backend");
if (removals.empty())
return {};
return {remove(removals.front().key, removals.front().expected_token, access)};
}

/// 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; }

virtual size_t conditionalBatchDeleteKeyLimit() const { return 1; }
virtual void pinConditionalBatchDeleteSupport(bool) {}

/// 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
Expand Up @@ -326,6 +326,20 @@ void InMemoryBackend::removeManyWriteOnce(const std::vector<WriteOnceKey> & keys
}
}

std::vector<Backend::RawRemoval> InMemoryBackend::removeManyByConditional(
const std::vector<RawConditionalRemove> & removals, TransportAccess &)
{
for (const RawConditionalRemove & removal : removals)
checkExpectedValue(removal.key, removal.expected_token);

std::lock_guard lock(mutex_);
std::vector<RawRemoval> results;
results.reserve(removals.size());
for (const RawConditionalRemove & removal : removals)
results.push_back(applyDelete(removal.key, removal.expected_token));
return results;
}

void InMemoryBackend::failNextBulkRemoveWith(std::exception_ptr error)
{
std::lock_guard lock(mutex_);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,8 @@ class InMemoryBackend : public Backend
/// `hold_deletes_` exactly as `remove` does: a held delete is queued and lands only on
/// `landPendingDelete`.
void removeManyWriteOnce(const std::vector<WriteOnceKey> & keys, TransportAccess & access) override;
std::vector<RawRemoval> removeManyByConditional(
const std::vector<RawConditionalRemove> & removals, TransportAccess & access) override;
/// The next `removeManyWriteOnce` throws `error` instead of deleting anything; one-shot, like
/// `failNextWriteWith`.
void failNextBulkRemoveWith(std::exception_ptr error);
Expand All @@ -62,6 +64,8 @@ class InMemoryBackend : public Backend

/// 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; }
size_t conditionalBatchDeleteKeyLimit() const override { return conditional_batch_delete_supported ? kBulkDeleteMaxKeys : 1; }
void pinConditionalBatchDeleteSupport(bool supported) override { conditional_batch_delete_supported = supported; }

/// 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
Expand Down Expand Up @@ -227,6 +231,7 @@ class InMemoryBackend : public Backend
std::vector<std::exception_ptr> armed_bulk_remove_failures_;
std::function<void()> before_bulk_remove_hook_;
size_t bulk_remove_calls_ = 0;
std::atomic_bool conditional_batch_delete_supported = false;
};

}
Original file line number Diff line number Diff line change
Expand Up @@ -154,4 +154,13 @@ void InstrumentedBackend::removeManyWriteOnce(const std::vector<WriteOnceKey> &
incrementCasEvent(classifyCasNs(key.str()), CasOp::Delete);
}

std::vector<Backend::RawRemoval> InstrumentedBackend::removeManyByConditional(
const std::vector<RawConditionalRemove> & removals, TransportAccess & access)
{
std::vector<RawRemoval> results = inner->removeManyByConditional(removals, access);
for (const RawConditionalRemove & removal : removals)
incrementCasEvent(classifyCasNs(removal.key), CasOp::Delete);
return results;
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,8 @@ class InstrumentedBackend final : public Backend
/// carried, the request counter says how many requests it took. Out of line, like `publish`, so
/// the header need not declare the `CASBulkDeleteRequests` extern.
void removeManyWriteOnce(const std::vector<WriteOnceKey> & keys, TransportAccess & access) override;
std::vector<RawRemoval> removeManyByConditional(
const std::vector<RawConditionalRemove> & removals, TransportAccess & access) override;

/// Count a create and a replacement separately, and each of them separately from its refusal:
/// they cost the same one request, but a pool whose creates are mostly refused and one whose
Expand Down Expand Up @@ -162,6 +164,8 @@ class InstrumentedBackend final : public Backend
uint64_t attemptTimeoutMs() const override { return inner->attemptTimeoutMs(); }
uint64_t attemptEnvelopeMs() const override { return inner->attemptEnvelopeMs(); }
size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); }
size_t conditionalBatchDeleteKeyLimit() const override { return inner->conditionalBatchDeleteKeyLimit(); }
void pinConditionalBatchDeleteSupport(bool supported) override { inner->pinConditionalBatchDeleteSupport(supported); }
bool refreshCredentials() override { return inner->refreshCredentials(); }

private:
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1009,6 +1009,18 @@ size_t ObjectStorageBackend::bulkDeleteKeyLimit() const
return mode == Mode::Native ? object_storage->batchDeleteKeyLimit() : kBulkDeleteMaxKeys;
}

size_t ObjectStorageBackend::conditionalBatchDeleteKeyLimit() const
{
if (mode != Mode::Native || !conditional_batch_delete_supported.load())
return 1;
return object_storage->batchDeleteKeyLimit();
}

void ObjectStorageBackend::pinConditionalBatchDeleteSupport(bool supported)
{
conditional_batch_delete_supported.store(supported);
}

void ObjectStorageBackend::removeManyWriteOnce(const std::vector<WriteOnceKey> & keys, TransportAccess & access)
{
if (keys.empty())
Expand Down Expand Up @@ -1036,6 +1048,48 @@ void ObjectStorageBackend::removeManyWriteOnce(const std::vector<WriteOnceKey> &
}
}

std::vector<Backend::RawRemoval> ObjectStorageBackend::removeManyByConditional(
const std::vector<RawConditionalRemove> & removals, TransportAccess & access)
{
if (removals.empty())
return {};
if (mode != Mode::Native)
{
std::vector<RawRemoval> results;
results.reserve(removals.size());
for (const RawConditionalRemove & removal : removals)
results.push_back(removeUnder(removal.key, removal.expected_token, controlRequest(access.attemptNo())));
return results;
}

std::vector<ConditionalRemoveObject> objects;
objects.reserve(removals.size());
for (const RawConditionalRemove & removal : removals)
{
if (!isValidTokenValue(dialect(), removal.expected_token))
throw Exception(ErrorCodes::LOGICAL_ERROR, "CAS backend: malformed conditional removal token for {}", removal.key);
objects.push_back({StoredObject(removal.key), removal.expected_token});
}
const auto storage_results = object_storage->removeObjectsIfTokensMatch(objects, controlRequest(access.attemptNo()));
if (storage_results.size() != removals.size())
throw Exception(ErrorCodes::LOGICAL_ERROR, "CAS backend: conditional batch removal returned {} results for {} keys", storage_results.size(), removals.size());

std::vector<RawRemoval> results;
results.reserve(storage_results.size());
for (const ConditionalRemoveResult & result : storage_results)
{
switch (result.outcome)
{
case ConditionalRemoveOutcome::Removed:
results.push_back(result.created_delete_marker ? RawRemoval::DeleteMarker : RawRemoval::Removed);
break;
case ConditionalRemoveOutcome::TokenMismatch: results.push_back(RawRemoval::Mismatch); break;
case ConditionalRemoveOutcome::NotFound: results.push_back(RawRemoval::Gone); break;
}
}
return results;
}

Backend::RawListPage ObjectStorageBackend::list(const String & prefix, const String & cursor, size_t limit, TransportAccess & access)
{
return listUnder(prefix, cursor, limit, controlRequest(access.attemptNo()));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -101,10 +101,14 @@ 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;
std::vector<RawRemoval> removeManyByConditional(
const std::vector<RawConditionalRemove> & removals, 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;
size_t conditionalBatchDeleteKeyLimit() const override;
void pinConditionalBatchDeleteSupport(bool supported) 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 Expand Up @@ -214,6 +218,7 @@ class ObjectStorageBackend final : public Backend
private:
const ObjectStoragePtr object_storage;
const Mode mode;
std::atomic_bool conditional_batch_delete_supported = false;
Dialect native_token_type = Dialect::ETag;
/// See the constructor: what the READ-class requests (read, head, list, remove) carry.
const bool single_attempt_control_plane;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,11 @@

#include <fmt/format.h>

#include <algorithm>
#include <array>
#include <ctime>
#include <limits>
#include <set>
#include <utility>

namespace ProfileEvents
Expand Down Expand Up @@ -613,6 +615,57 @@ void CasOperation::removeManyWriteOnce(const std::vector<WriteOnceKey> & keys, c
});
}

std::vector<Removal> CasOperation::removeManyByConditional(
const std::vector<ConditionalRemove> & removals, const Retry & policy)
{
const size_t limit = std::min(kBulkDeleteMaxKeys, owner.backend->conditionalBatchDeleteKeyLimit());
if (removals.size() > limit)
throw Exception(ErrorCodes::LOGICAL_ERROR,
"CAS removeManyByConditional: {} keys in one chunk, the limit is {}; the consumer chunks its input",
removals.size(), limit);
if (removals.empty())
return {};

std::set<String> keys;
std::vector<Backend::RawConditionalRemove> raw_removals;
raw_removals.reserve(removals.size());
for (const ConditionalRemove & removal : removals)
{
if (!keys.emplace(removal.key).second)
throw Exception(ErrorCodes::LOGICAL_ERROR,
"CAS removeManyByConditional: duplicate key '{}' in one chunk", removal.key);
raw_removals.push_back({removal.key, owner.valueFor(removal.key, removal.expected_token)});
}

const Retry::Bound bound = policy.bind(owner.now_ms());
const String subject = fmt::format("{} (+{} keys)", removals.front().key, removals.size() - 1);
const std::vector<Backend::RawRemoval> raw_results = readLoop(
"removeManyByConditional", subject, policy, bound, [&](auto & access)
{
return owner.backend->removeManyByConditional(raw_removals, access);
});
if (raw_results.size() != removals.size())
throw Exception(ErrorCodes::LOGICAL_ERROR,
"CAS removeManyByConditional: backend returned {} results for {} keys", raw_results.size(), removals.size());

std::vector<Removal> results;
results.reserve(raw_results.size());
for (const Backend::RawRemoval raw : raw_results)
{
switch (raw)
{
case Backend::RawRemoval::Removed: results.push_back(Removal::Removed); break;
case Backend::RawRemoval::Gone: results.push_back(Removal::Gone); break;
case Backend::RawRemoval::Mismatch: results.push_back(Removal::Mismatch); break;
case Backend::RawRemoval::DeleteMarker:
throw Exception(ErrorCodes::CAS_DELETE_MARKER,
"CAS conditional batch removal archived a noncurrent version instead of reclaiming an object: "
"the bucket has object versioning enabled");
}
}
return results;
}

Removal CasOperation::remove(const String & key, const Etag & seen, const Retry & policy)
{
return removeUnder(key, owner.valueFor(key, seen), policy, policy.bind(owner.now_ms()));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -234,6 +234,12 @@ class CasRequests
/// The cap on one `removeManyWriteOnce` chunk -- also the ceiling a batch-delete request can carry.
inline constexpr size_t kBulkDeleteMaxKeys = 1000;

struct ConditionalRemove
{
String key;
Etag expected_token;
};

/// One admitted operation: the unit a policy, a fence generation and a liveness predicate apply to.
/// Move-only, and every request it makes re-checks its admission -- before each attempt, before each
/// sleep, and once more after a proven commit, so a write whose fence was lost while it was in flight
Expand Down Expand Up @@ -288,6 +294,7 @@ class CasOperation
/// success. Throws when the policy is exhausted. More keys than the cap is a caller bug: the
/// consumer chunks, so that every chunk that succeeded is recorded before a later one can fail.
void removeManyWriteOnce(const std::vector<WriteOnceKey> & keys, const Retry & policy);
std::vector<Removal> removeManyByConditional(const std::vector<ConditionalRemove> & removals, const Retry & policy);
/// The one primitive that reports failure as a value, so the policy reissues on the OUTCOME:
/// `Indeterminate` is retried, the four authoritative outcomes return at once, and an
/// `Indeterminate` that outlives the bound is returned rather than thrown. Admission refused before
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,14 @@ class ThrottlingBackend final : public Backend
inner->removeManyWriteOnce(keys, access);
}

std::vector<RawRemoval> removeManyByConditional(
const std::vector<RawConditionalRemove> & removals, TransportAccess & access) override
{
for (const RawConditionalRemove & removal : removals)
refuseOrPass(removal.key);
return inner->removeManyByConditional(removals, access);
}

std::expected<String, RawConflict> write(const String & key, const String & bytes,
const std::optional<String> & expected_value, TransportAccess & access) override
{
Expand Down Expand Up @@ -126,6 +134,8 @@ class ThrottlingBackend final : public Backend
uint64_t attemptTimeoutMs() const override { return inner->attemptTimeoutMs(); }
uint64_t attemptEnvelopeMs() const override { return inner->attemptEnvelopeMs(); }
size_t bulkDeleteKeyLimit() const override { return inner->bulkDeleteKeyLimit(); }
size_t conditionalBatchDeleteKeyLimit() const override { return inner->conditionalBatchDeleteKeyLimit(); }
void pinConditionalBatchDeleteSupport(bool supported) override { inner->pinConditionalBatchDeleteSupport(supported); }
bool refreshCredentials() override { return inner->refreshCredentials(); }
void checkPoolPreconditions() override { inner->checkPoolPreconditions(); }
void checkSkipAccessCheckSupport() override { inner->checkSkipAccessCheckSupport(); }
Expand Down
Loading
Loading