Copilot commented on code in PR #3502:
URL: https://github.com/apache/kvrocks/pull/3502#discussion_r3310098250
##########
src/types/redis_stream.cc:
##########
@@ -394,6 +395,337 @@ rocksdb::Status Stream::DeletePelEntries(engine::Context
&ctx, const Slice &stre
return storage_->Write(ctx, storage_->DefaultWriteOptions(),
batch->GetWriteBatch());
}
+rocksdb::Status Stream::getGroupNames(engine::Context &ctx, const std::string
&ns_key, const StreamMetadata &metadata,
+ std::vector<std::string> *group_names) {
+ group_names->clear();
+
+ std::string subkey_type_delimiter;
+ PutFixed64(&subkey_type_delimiter, UINT64_MAX);
+ PutFixed8(&subkey_type_delimiter,
static_cast<uint8_t>(StreamSubkeyType::StreamConsumerGroupMetadata));
+
+ std::string next_version_prefix_key =
+ InternalKey(ns_key, subkey_type_delimiter, metadata.version + 1,
storage_->IsSlotIdEncoded()).Encode();
+ std::string prefix_key =
+ InternalKey(ns_key, subkey_type_delimiter, metadata.version,
storage_->IsSlotIdEncoded()).Encode();
+
+ rocksdb::ReadOptions read_options = ctx.DefaultScanOptions();
+ rocksdb::Slice upper_bound(next_version_prefix_key);
+ read_options.iterate_upper_bound = &upper_bound;
+ rocksdb::Slice lower_bound(prefix_key);
+ read_options.iterate_lower_bound = &lower_bound;
+
+ auto iter = util::UniqueIterator(ctx, read_options, stream_cf_handle_);
+ for (iter->SeekToFirst(); iter->Valid(); iter->Next()) {
+ if (identifySubkeyType(iter->key()) !=
StreamSubkeyType::StreamConsumerGroupMetadata) {
+ continue;
+ }
+ group_names->push_back(groupNameFromInternalKey(iter->key()));
+ }
+ return iter->status();
+}
+
+rocksdb::Status Stream::deleteEntryAndUpdateMeta(rocksdb::WriteBatchBase
*batch, const std::string &entry_key,
+ const StreamEntryID &id,
StreamMetadata *metadata,
+ uint64_t *deleted_cnt) {
+ rocksdb::Status s = batch->Delete(stream_cf_handle_, entry_key);
+ if (!s.ok()) return s;
+ (*deleted_cnt)++;
+
+ if (metadata->max_deleted_entry_id < id) {
+ metadata->max_deleted_entry_id = id;
+ }
+
+ return rocksdb::Status::OK();
+}
+
+rocksdb::Status Stream::processOneEntryDeletion(
+ engine::Context &ctx, const rocksdb::ReadOptions &read_options, const
std::string &ns_key, StreamMetadata *metadata,
+ const StreamEntryID &id, const std::vector<std::string> &all_groups,
StreamDeleteOption option,
+ rocksdb::WriteBatchBase *batch, bool *batch_modified, uint64_t
*deleted_cnt, StreamEntryDeleteResult *result,
+ std::map<std::string, uint64_t> *group_pending_decrements,
+ std::map<std::string, std::map<std::string, uint64_t>>
*consumer_pending_decrements) {
+ std::string entry_key = internalKeyFromEntryID(ns_key, *metadata, id);
+ std::string value;
+ auto s = storage_->Get(ctx, read_options, stream_cf_handle_, entry_key,
&value);
+ if (!s.ok() && !s.IsNotFound()) {
+ return s;
+ }
+
+ if (s.IsNotFound()) {
+ if (option == StreamDeleteOption::DelRef) {
+ s = cleanPelFromAllGroups(ctx, ns_key, *metadata, id, batch,
batch_modified, all_groups, group_pending_decrements,
+ consumer_pending_decrements);
+ if (!s.ok()) return s;
+ }
+ *result = StreamEntryDeleteResult::kEntryNotFound;
+ return rocksdb::Status::OK();
+ }
+
+ switch (option) {
+ case StreamDeleteOption::KeepRef:
+ s = deleteEntryAndUpdateMeta(batch, entry_key, id, metadata,
deleted_cnt);
+ if (!s.ok()) return s;
+ *batch_modified = true;
+ *result = StreamEntryDeleteResult::kEntryDeleted;
+ break;
+
+ case StreamDeleteOption::DelRef:
+ s = deleteEntryAndUpdateMeta(batch, entry_key, id, metadata,
deleted_cnt);
+ if (!s.ok()) return s;
+ s = cleanPelFromAllGroups(ctx, ns_key, *metadata, id, batch,
batch_modified, all_groups, group_pending_decrements,
+ consumer_pending_decrements);
+ if (!s.ok()) return s;
+ *batch_modified = true;
+ *result = StreamEntryDeleteResult::kEntryDeleted;
+ break;
+
+ case StreamDeleteOption::Acked: {
+ if (all_groups.empty()) {
+ *result = StreamEntryDeleteResult::kEntrySkipped;
+ return rocksdb::Status::OK();
+ }
+ bool all_acked = false;
+ s = isAckedByAllGroups(ctx, ns_key, *metadata, id, all_groups,
&all_acked);
+ if (!s.ok()) return s;
+ if (all_acked) {
+ s = deleteEntryAndUpdateMeta(batch, entry_key, id, metadata,
deleted_cnt);
+ if (!s.ok()) return s;
+ *batch_modified = true;
+ *result = StreamEntryDeleteResult::kEntryDeleted;
+ } else {
+ *result = StreamEntryDeleteResult::kEntrySkipped;
+ }
+ break;
+ }
+ }
+
+ return rocksdb::Status::OK();
+}
+
+rocksdb::Status Stream::cleanPelFromAllGroups(
+ engine::Context &ctx, const std::string &ns_key, const StreamMetadata
&metadata, const StreamEntryID &id,
+ rocksdb::WriteBatchBase *batch, bool *batch_modified, const
std::vector<std::string> &group_names,
+ std::map<std::string, uint64_t> *group_pending_decrements,
+ std::map<std::string, std::map<std::string, uint64_t>>
*consumer_pending_decrements) {
+ for (const auto &group_name : group_names) {
+ std::string pel_key = internalPelKeyFromGroupAndEntryId(ns_key, metadata,
group_name, id);
+ std::string pel_value;
+ auto pel_s = storage_->Get(ctx, ctx.GetReadOptions(), stream_cf_handle_,
pel_key, &pel_value);
+ if (pel_s.ok()) {
+ rocksdb::Status s = batch->Delete(stream_cf_handle_, pel_key);
+ if (!s.ok()) return s;
+ *batch_modified = true;
+
+ auto pel_entry = decodeStreamPelEntryValue(pel_value);
+ (*group_pending_decrements)[group_name]++;
+ (*consumer_pending_decrements)[group_name][pel_entry.consumer_name]++;
+ } else if (!pel_s.IsNotFound()) {
+ return pel_s;
+ }
+ }
+ return rocksdb::Status::OK();
+}
+
+rocksdb::Status Stream::flushPendingNumberUpdates(
+ engine::Context &ctx, const std::string &ns_key, const StreamMetadata
&metadata, rocksdb::WriteBatchBase *batch,
+ const std::map<std::string, uint64_t> &group_pending_decrements,
+ const std::map<std::string, std::map<std::string, uint64_t>>
&consumer_pending_decrements) {
+ for (const auto &[group_name, decrement] : group_pending_decrements) {
+ auto group_key = internalKeyFromGroupName(ns_key, metadata, group_name);
+ std::string group_value;
+ auto s = storage_->Get(ctx, ctx.GetReadOptions(), stream_cf_handle_,
group_key, &group_value);
+ if (!s.ok() && !s.IsNotFound()) {
+ return s;
+ }
+ if (s.ok()) {
+ auto group_meta = decodeStreamConsumerGroupMetadataValue(group_value);
+ group_meta.pending_number -= decrement;
+ s = batch->Put(stream_cf_handle_, group_key,
encodeStreamConsumerGroupMetadataValue(group_meta));
+ if (!s.ok()) return s;
+ }
+ }
+
+ for (const auto &[group_name, consumers] : consumer_pending_decrements) {
+ for (const auto &[consumer_name, decrement] : consumers) {
+ auto consumer_key = internalKeyFromConsumerName(ns_key, metadata,
group_name, consumer_name);
+ std::string consumer_value;
+ auto s = storage_->Get(ctx, ctx.GetReadOptions(), stream_cf_handle_,
consumer_key, &consumer_value);
+ if (!s.ok() && !s.IsNotFound()) {
+ return s;
+ }
+ if (s.ok()) {
+ auto consumer_meta = decodeStreamConsumerMetadataValue(consumer_value);
+ consumer_meta.pending_number -= decrement;
+ s = batch->Put(stream_cf_handle_, consumer_key,
encodeStreamConsumerMetadataValue(consumer_meta));
+ if (!s.ok()) return s;
+ }
+ }
+ }
+
+ return rocksdb::Status::OK();
+}
+
+rocksdb::Status Stream::isAckedByAllGroups(engine::Context &ctx, const
std::string &ns_key,
+ const StreamMetadata &metadata,
const StreamEntryID &id,
+ const std::vector<std::string>
&all_groups, bool *all_acked) {
+ *all_acked = true;
+ for (const auto &group_name : all_groups) {
+ std::string pel_key = internalPelKeyFromGroupAndEntryId(ns_key, metadata,
group_name, id);
+ std::string pel_value;
+ auto pel_s = storage_->Get(ctx, ctx.GetReadOptions(), stream_cf_handle_,
pel_key, &pel_value);
+ if (pel_s.ok()) {
+ *all_acked = false;
+ return rocksdb::Status::OK();
+ }
+ if (!pel_s.IsNotFound()) {
+ return pel_s;
+ }
+ }
+ return rocksdb::Status::OK();
+}
+
+rocksdb::Status Stream::DeleteEntriesWithOption(engine::Context &ctx, const
Slice &stream_name,
+ const
std::vector<StreamEntryID> &ids, StreamDeleteOption option,
+ std::vector<int> *results) {
+ results->assign(ids.size(),
static_cast<int>(StreamEntryDeleteResult::kEntryNotFound));
+
+ std::string ns_key = AppendNamespacePrefix(stream_name);
+
+ StreamMetadata metadata(false);
+ rocksdb::Status s = GetMetadata(ctx, ns_key, &metadata);
+ if (!s.ok()) {
+ return s.IsNotFound() ? rocksdb::Status::OK() : s;
+ }
+
+ if (ids.empty()) {
+ return rocksdb::Status::OK();
+ }
+
+ auto batch = storage_->GetWriteBatchBase();
+ std::string option_str;
+ switch (option) {
+ case StreamDeleteOption::DelRef:
+ option_str = "DELREF";
+ break;
+ case StreamDeleteOption::Acked:
+ option_str = "ACKED";
+ break;
+ default:
+ option_str = "KEEPREF";
+ break;
+ }
+ WriteBatchLogData log_data(kRedisStream, {"XDELEX", option_str});
+ s = batch->PutLogData(log_data.Encode());
+ if (!s.ok()) return s;
+
+ std::string next_version_prefix_key =
+ InternalKey(ns_key, "", metadata.version + 1,
storage_->IsSlotIdEncoded()).Encode();
+ std::string prefix_key = InternalKey(ns_key, "", metadata.version,
storage_->IsSlotIdEncoded()).Encode();
+
+ rocksdb::ReadOptions read_options = ctx.DefaultScanOptions();
+ rocksdb::Slice upper_bound(next_version_prefix_key);
+ read_options.iterate_upper_bound = &upper_bound;
+ rocksdb::Slice lower_bound(prefix_key);
+ read_options.iterate_lower_bound = &lower_bound;
+
+ auto iter = util::UniqueIterator(ctx, read_options, stream_cf_handle_);
+
+ std::vector<std::string> all_groups;
+ bool need_groups = (option == StreamDeleteOption::DelRef || option ==
StreamDeleteOption::Acked);
+ if (need_groups) {
+ s = getGroupNames(ctx, ns_key, metadata, &all_groups);
+ if (!s.ok()) return s;
+ }
+
+ uint64_t deleted_cnt = 0;
+ bool batch_modified = false;
+ std::map<std::string, uint64_t> group_pending_decrements;
+ std::map<std::string, std::map<std::string, uint64_t>>
consumer_pending_decrements;
+
+ std::unordered_set<std::string> seen_keys;
+ seen_keys.reserve(ids.size());
+ std::unordered_set<std::string> deleted_entry_keys;
+ deleted_entry_keys.reserve(ids.size());
+ StreamEntryID orig_first_entry_id = metadata.first_entry_id;
+ StreamEntryID orig_last_entry_id = metadata.last_entry_id;
+
+ for (size_t i = 0; i < ids.size(); i++) {
+ const auto &id = ids[i];
+
+ std::string entry_key = internalKeyFromEntryID(ns_key, metadata, id);
+
+ if (!seen_keys.insert(entry_key).second) {
+ (*results)[i] =
static_cast<int>(StreamEntryDeleteResult::kEntryNotFound);
+ continue;
+ }
+
+ StreamEntryDeleteResult result = StreamEntryDeleteResult::kEntryNotFound;
+ s = processOneEntryDeletion(ctx, read_options, ns_key, &metadata, id,
all_groups, option, batch.Get(),
+ &batch_modified, &deleted_cnt, &result,
&group_pending_decrements,
+ &consumer_pending_decrements);
+ if (!s.ok()) return s;
+ if (result == StreamEntryDeleteResult::kEntryDeleted) {
+ deleted_entry_keys.insert(entry_key);
+ }
+ (*results)[i] = static_cast<int>(result);
+ }
+
+ if (deleted_cnt > 0) {
+ metadata.size -= deleted_cnt;
+
+ if (metadata.size == 0) {
+ metadata.first_entry_id.Clear();
+ metadata.last_entry_id.Clear();
+ metadata.recorded_first_entry_id.Clear();
+ } else {
+ bool first_deleted =
deleted_entry_keys.count(internalKeyFromEntryID(ns_key, metadata,
orig_first_entry_id)) > 0;
+ bool last_deleted =
deleted_entry_keys.count(internalKeyFromEntryID(ns_key, metadata,
orig_last_entry_id)) > 0;
+
+ if (first_deleted) {
+ iter->SeekToFirst();
+ while (iter->Valid() &&
deleted_entry_keys.count(iter->key().ToString()) > 0) {
+ iter->Next();
+ }
+ if (iter->Valid()) {
+ metadata.first_entry_id = entryIDFromInternalKey(iter->key());
+ metadata.recorded_first_entry_id = metadata.first_entry_id;
+ } else {
+ metadata.first_entry_id.Clear();
+ metadata.recorded_first_entry_id.Clear();
+ }
+ }
+ if (last_deleted) {
+ iter->SeekToLast();
+ while (iter->Valid() &&
deleted_entry_keys.count(iter->key().ToString()) > 0) {
+ iter->Prev();
+ }
+ if (iter->Valid()) {
+ metadata.last_entry_id = entryIDFromInternalKey(iter->key());
+ } else {
+ metadata.last_entry_id.Clear();
+ }
+ }
Review Comment:
When recalculating `metadata.last_entry_id` after deletions, the iterator is
positioned with `SeekToLast()` over the entire stream keyspace for the metadata
version. Stream group/consumer/PEL keys sort after stream entry keys (they
start with the UINT64_MAX delimiter), so `SeekToLast()` will often land on
non-entry keys and `entryIDFromInternalKey()` will parse garbage, corrupting
the stream metadata. The boundary scan should skip non-`StreamEntry` keys (and
also skip keys deleted in this batch).
##########
src/commands/cmd_stream.cc:
##########
@@ -268,6 +268,73 @@ class CommandXDel : public Commander {
std::vector<redis::StreamEntryID> ids_;
};
+class CommandXDelEx : public Commander {
+ public:
+ Status Parse(const std::vector<std::string> &args) override {
+ CommandParser parser(args, 1);
+ stream_name_ = GET_OR_RET(parser.TakeStr());
+
+ option_ = redis::StreamDeleteOption::KeepRef;
+
+ while (parser.Good() && !util::EqualICase(parser.RawPeek(), "IDS")) {
+ if (parser.EatEqICase("KEEPREF")) {
+ option_ = redis::StreamDeleteOption::KeepRef;
+ } else if (parser.EatEqICase("DELREF")) {
+ option_ = redis::StreamDeleteOption::DelRef;
+ } else if (parser.EatEqICase("ACKED")) {
+ option_ = redis::StreamDeleteOption::Acked;
+ } else {
+ return parser.InvalidSyntax();
+ }
+ }
+
+ if (!parser.EatEqICase("IDS")) {
+ return {Status::RedisParseErr, "syntax error, expected IDS keyword"};
+ }
+
+ auto numids_result = parser.TakeInt<int64_t>();
+ if (!numids_result.IsOK()) {
+ return {Status::RedisParseErr, errValueNotInteger};
+ }
+ int64_t numids = numids_result.GetValue();
+ if (numids <= 0) {
+ return {Status::RedisParseErr, "numids must be positive"};
+ }
+
+ for (int64_t i = 0; i < numids; i++) {
+ auto id_str = GET_OR_RET(parser.TakeStr());
+ redis::StreamEntryID id;
+ auto s = ParseStreamEntryID(id_str, &id);
+ if (!s.IsOK()) return s;
+ entry_ids_.emplace_back(id);
+ }
+
+ return Status::OK();
Review Comment:
`XDELEX` parsing currently ignores any trailing arguments after the declared
`numids` IDs (e.g. `XDELEX s IDS 1 1-0 EXTRA` would be accepted). Most commands
in this codebase treat leftover tokens as a syntax error; accepting extras can
hide client bugs and makes the command grammar ambiguous.
##########
src/types/redis_stream.h:
##########
@@ -22,10 +22,12 @@
#include <rocksdb/status.h>
+#include <map>
#include <optional>
#include <string>
#include <vector>
+#include "common/db_util.h"
#include "storage/redis_db.h"
#include "storage/redis_metadata.h"
Review Comment:
`src/types/redis_stream.h` now includes `common/db_util.h`, but this header
doesn't use any symbols from it. Keeping it here increases compile-time
coupling (it pulls in RocksDB DB headers, storage, etc.). Prefer removing it
and keep `db_util.h` only in the .cc where `util::UniqueIterator` is actually
used.
--
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]