Copilot commented on code in PR #3504:
URL: https://github.com/apache/kvrocks/pull/3504#discussion_r3310099980
##########
src/types/redis_stream.cc:
##########
@@ -394,6 +395,431 @@ 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();
+ } else if (!pel_s.IsNotFound()) {
+ return pel_s;
+ }
+ }
+ return rocksdb::Status::OK();
+}
+
+rocksdb::Status Stream::DeleteEntriesAndAck(engine::Context &ctx, const Slice
&stream_name,
+ const std::string &group_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::NotFound("NOGROUP No such
consumer group '" + group_name + "'") : s;
+ }
+
+ if (ids.empty()) {
+ return rocksdb::Status::OK();
+ }
+
+ std::string group_key = internalKeyFromGroupName(ns_key, metadata,
group_name);
+ std::string get_group_value;
+ s = storage_->Get(ctx, ctx.GetReadOptions(), stream_cf_handle_, group_key,
&get_group_value);
+ if (!s.ok()) {
+ if (s.IsNotFound()) {
+ return rocksdb::Status::NotFound("NOGROUP No such consumer group '" +
group_name + "'");
+ }
+ return s;
+ }
+
+ 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, {"XACKDEL", group_name,
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;
+ }
+
+ std::map<std::string, uint64_t> consumer_acknowledges;
+ std::map<std::string, uint64_t> other_group_pending_decrements;
+ std::map<std::string, std::map<std::string, uint64_t>>
other_consumer_pending_decrements;
+ uint64_t deleted_cnt = 0;
+ uint64_t acknowledged_cnt = 0;
+ bool batch_modified = false;
+
+ 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;
+ }
+
+ std::string value;
+ s = storage_->Get(ctx, read_options, stream_cf_handle_, entry_key, &value);
+ if (!s.ok() && !s.IsNotFound()) {
+ return s;
+ }
+ bool stream_entry_exists = s.ok();
+
+ std::string pel_key = internalPelKeyFromGroupAndEntryId(ns_key, metadata,
group_name, id);
+ std::string pel_value;
+ s = storage_->Get(ctx, ctx.GetReadOptions(), stream_cf_handle_, pel_key,
&pel_value);
+ if (!s.ok() && !s.IsNotFound()) {
+ return s;
+ }
+ bool pel_entry_exists = s.ok();
+
+ if (!stream_entry_exists && !pel_entry_exists) {
+ (*results)[i] =
static_cast<int>(StreamEntryDeleteResult::kEntryNotFound);
+ continue;
+ }
+
+ bool pel_cleaned = false;
+ if (pel_entry_exists) {
+ s = batch->Delete(stream_cf_handle_, pel_key);
+ if (!s.ok()) return s;
+ pel_cleaned = true;
+ acknowledged_cnt++;
+ batch_modified = true;
+
+ auto pel_entry = decodeStreamPelEntryValue(pel_value);
+ consumer_acknowledges[pel_entry.consumer_name]++;
+ }
Review Comment:
The implementation always deletes the PEL entry for the specified
`group_name` when it exists (and decrements pending counters) before applying
the `KEEPREF/DELREF/ACKED` strategy. This appears to conflict with the PR
description where `KEEPREF` says it should not touch PEL references, and it
also makes requests return `1` when only the group's PEL existed but the stream
entry was already missing (via the `pel_cleaned` fallback). Please align the
behavior and return codes with the documented semantics (either update the code
to only touch PELs when requested, or update the command spec/tests to match
the actual behavior).
##########
src/types/redis_stream.cc:
##########
@@ -394,6 +395,431 @@ 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();
+ } else if (!pel_s.IsNotFound()) {
+ return pel_s;
+ }
+ }
+ return rocksdb::Status::OK();
+}
+
+rocksdb::Status Stream::DeleteEntriesAndAck(engine::Context &ctx, const Slice
&stream_name,
+ const std::string &group_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::NotFound("NOGROUP No such
consumer group '" + group_name + "'") : s;
+ }
+
+ if (ids.empty()) {
+ return rocksdb::Status::OK();
+ }
+
+ std::string group_key = internalKeyFromGroupName(ns_key, metadata,
group_name);
+ std::string get_group_value;
+ s = storage_->Get(ctx, ctx.GetReadOptions(), stream_cf_handle_, group_key,
&get_group_value);
+ if (!s.ok()) {
+ if (s.IsNotFound()) {
+ return rocksdb::Status::NotFound("NOGROUP No such consumer group '" +
group_name + "'");
+ }
+ return s;
+ }
+
+ 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, {"XACKDEL", group_name,
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;
+ }
+
+ std::map<std::string, uint64_t> consumer_acknowledges;
+ std::map<std::string, uint64_t> other_group_pending_decrements;
+ std::map<std::string, std::map<std::string, uint64_t>>
other_consumer_pending_decrements;
+ uint64_t deleted_cnt = 0;
+ uint64_t acknowledged_cnt = 0;
+ bool batch_modified = false;
+
+ 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;
+ }
+
+ std::string value;
+ s = storage_->Get(ctx, read_options, stream_cf_handle_, entry_key, &value);
+ if (!s.ok() && !s.IsNotFound()) {
+ return s;
+ }
+ bool stream_entry_exists = s.ok();
+
+ std::string pel_key = internalPelKeyFromGroupAndEntryId(ns_key, metadata,
group_name, id);
+ std::string pel_value;
+ s = storage_->Get(ctx, ctx.GetReadOptions(), stream_cf_handle_, pel_key,
&pel_value);
+ if (!s.ok() && !s.IsNotFound()) {
+ return s;
+ }
+ bool pel_entry_exists = s.ok();
+
+ if (!stream_entry_exists && !pel_entry_exists) {
+ (*results)[i] =
static_cast<int>(StreamEntryDeleteResult::kEntryNotFound);
+ continue;
+ }
+
+ bool pel_cleaned = false;
+ if (pel_entry_exists) {
+ s = batch->Delete(stream_cf_handle_, pel_key);
+ if (!s.ok()) return s;
+ pel_cleaned = true;
+ acknowledged_cnt++;
+ batch_modified = true;
+
+ auto pel_entry = decodeStreamPelEntryValue(pel_value);
+ consumer_acknowledges[pel_entry.consumer_name]++;
+ }
+
+ const std::vector<std::string> *groups_for_cleanup = &all_groups;
+ std::vector<std::string> filtered_groups;
+ if ((option == StreamDeleteOption::DelRef || option ==
StreamDeleteOption::Acked) && pel_entry_exists) {
+ for (const auto &g : all_groups) {
+ if (g != group_name) {
+ filtered_groups.push_back(g);
+ }
+ }
+ groups_for_cleanup = &filtered_groups;
+ }
+
+ StreamEntryDeleteResult del_result =
StreamEntryDeleteResult::kEntryNotFound;
+ s = processOneEntryDeletion(ctx, read_options, ns_key, &metadata, id,
*groups_for_cleanup, option, batch.Get(),
+ &batch_modified, &deleted_cnt, &del_result,
&other_group_pending_decrements,
+ &other_consumer_pending_decrements);
+ if (!s.ok()) return s;
+
+ if (del_result == StreamEntryDeleteResult::kEntrySkipped && option ==
StreamDeleteOption::Acked &&
+ pel_entry_exists && filtered_groups.empty() && !all_groups.empty()) {
+ s = deleteEntryAndUpdateMeta(batch.Get(), entry_key, id, &metadata,
&deleted_cnt);
+ if (!s.ok()) return s;
+ batch_modified = true;
+ del_result = StreamEntryDeleteResult::kEntryDeleted;
+ }
+
+ if (del_result == StreamEntryDeleteResult::kEntryDeleted) {
+ deleted_entry_keys.insert(entry_key);
+ }
+
+ if (del_result != StreamEntryDeleteResult::kEntryNotFound) {
+ (*results)[i] = static_cast<int>(del_result);
+ } else if (pel_cleaned) {
+ (*results)[i] = static_cast<int>(StreamEntryDeleteResult::kEntryDeleted);
+ }
+ }
+
+ if (deleted_cnt > 0 || acknowledged_cnt > 0) {
+ 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:
`DeleteEntriesAndAck` recalculates `metadata.last_entry_id` using
`iter->SeekToLast()` over a scan range that includes non-entry stream subkeys
(consumer/group metadata, PEL). The last key in this range will typically be a
non-entry key (it starts with `UINT64_MAX`), so
`entryIDFromInternalKey(iter->key())` can decode garbage and corrupt
`last_entry_id`/subsequent stream behavior. Restrict the iterator/range to
`StreamSubkeyType::StreamEntry` only, or skip backwards until
`identifySubkeyType(key)==StreamEntry` (and not in `deleted_entry_keys`) before
decoding the ID.
##########
src/storage/batch_extractor.cc:
##########
@@ -397,8 +405,25 @@ rocksdb::Status WriteBatchExtractor::DeleteCF(uint32_t
column_family_id, const S
Slice encoded_id = ikey.GetSubKey();
redis::StreamEntryID entry_id;
GetFixed64(&encoded_id, &entry_id.ms);
+
+ if (entry_id.ms == UINT64_MAX) {
+ return rocksdb::Status::OK();
+ }
+
GetFixed64(&encoded_id, &entry_id.seq);
- command_args = {"XDEL", ikey.GetKey().ToString(), entry_id.ToString()};
+ std::string entry_id_str = entry_id.ToString();
+ std::string user_key = ikey.GetKey().ToString();
+
+ auto args = log_data_.GetArguments();
+ if (!args->empty()) {
+ if ((*args)[0] == "XACKDEL" && args->size() >= 3) {
+ command_args = {(*args)[0], user_key, (*args)[1], (*args)[2], "IDS",
"1", entry_id_str};
+ } else {
+ command_args = {"XDEL", user_key, entry_id_str};
+ }
+ } else {
+ command_args = {"XDEL", user_key, entry_id_str};
+ }
}
Review Comment:
In `DeleteCF` for the Stream column family, `ns` is never populated from
`InternalKey`, so extracted `XDEL`/`XACKDEL` commands are appended under an
empty namespace. The Stream CF path also lacks the `slot_range_` key-slot
filtering that other CF branches apply, which can cause slot migration to emit
commands for keys outside the requested slot range.
--
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]