Copilot commented on code in PR #3502:
URL: https://github.com/apache/kvrocks/pull/3502#discussion_r3367465865


##########
src/storage/batch_extractor.cc:
##########
@@ -38,6 +41,8 @@ void WriteBatchExtractor::LogData(const rocksdb::Slice &blob) 
{
     // Redis type log data
     if (auto s = log_data_.Decode(blob); !s.IsOK()) {
       WARN("Failed to decode Redis type log: {}", s.Msg());
+    } else {
+      seen_xdelex_entry_keys_.clear();
     }
   }
 }

Review Comment:
   `seen_xdelex_entry_keys_` is cleared only when `log_data_.Decode(blob)` 
succeeds. If decoding fails (or a different log blob type is encountered), the 
dedup set can leak across batches and suppress emitting later XDELEX/XDEL 
commands for the same (namespace,key,id). Clear the set unconditionally at the 
start of `LogData()` so dedup is always scoped to the current batch.



##########
src/types/redis_stream.cc:
##########
@@ -394,6 +395,349 @@ 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()) {
+    // Group metadata keys sort before consumer/PEL keys, so no more groups 
once we see a different type.
+    if (identifySubkeyType(iter->key()) != 
StreamSubkeyType::StreamConsumerGroupMetadata) {
+      break;
+    }
+    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::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));

Review Comment:
   `pending_number` is a `uint64_t`; subtracting `decrement` without checking 
can underflow and wrap to a huge value if consumer metadata is out-of-sync. 
Clamp to 0 (or return corruption) when `pending_number < decrement`.



##########
src/types/redis_stream.cc:
##########
@@ -394,6 +395,349 @@ 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()) {
+    // Group metadata keys sort before consumer/PEL keys, so no more groups 
once we see a different type.
+    if (identifySubkeyType(iter->key()) != 
StreamSubkeyType::StreamConsumerGroupMetadata) {
+      break;
+    }
+    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::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));

Review Comment:
   `pending_number` is a `uint64_t`; subtracting `decrement` without checking 
can underflow and wrap to a huge value if metadata is out-of-sync (e.g., a PEL 
key exists but pending_number is already 0). Clamp to 0 (or return corruption) 
when `pending_number < decrement` to avoid invalid pending counts.



-- 
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