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]

Reply via email to