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


##########
src/commands/cmd_tdigest.cc:
##########
@@ -556,6 +556,57 @@ class CommandTDigestTrimmedMean : public Commander {
   double high_cut_quantile_;
 };
 
+class CommandTDigestCDF : public Commander {
+  Status Parse(const std::vector<std::string> &args) override {
+    if (args.size() == 2) return {Status::RedisParseErr, 
errWrongNumOfArguments};
+    key_name_ = args[1];
+    values_.reserve(args.size() - 2);
+    for (size_t i = 2; i < args.size(); i++) {
+      auto value = ParseFloat(args[i]);
+      if (!value) {
+        return {Status::RedisParseErr, errValueIsNotFloat};
+      }
+      values_.push_back(*value);
+    }
+    return Status::OK();
+  }
+
+  Status Execute(engine::Context &ctx, Server *srv, Connection *conn, 
std::string *output) override {
+    TDigest tdigest(srv->storage, conn->GetNamespace());
+    std::vector<std::string> cdf_result;
+    TDigestCDFResult result;
+    TDigestMetadata metadata;
+    auto meta_status = tdigest.GetMetaData(ctx, key_name_, &metadata);
+    std::vector<std::string> nan_results(values_.size(), "nan");
+    if (!meta_status.ok()) {
+      if (meta_status.IsNotFound()) {
+        return {Status::RedisExecErr, errKeyNotFound};
+      }
+      *output = redis::MultiBulkString(RESP::v2, nan_results);
+      return Status::OK();
+    }

Review Comment:
   `TDIGEST.CDF` currently swallows unexpected metadata read errors by 
returning an OK reply with `nan` values. This can hide real storage/corruption 
issues from clients; other TDigest commands propagate non-NotFound errors as 
`RedisExecErr` (e.g. `tdigest.info`, `tdigest.quantile`).



##########
tests/cppunit/types/tdigest_test.cc:
##########
@@ -948,3 +949,136 @@ TEST_F(RedisTDigestTest, 
MergeWithUserSpecifiedCompression) {
   // Verify total observations: dest(1) + src(1) = 2
   EXPECT_EQ(metadata.total_observations, 2);
 }
+
+TEST_F(RedisTDigestTest, CDF_Test) {
+  std::string cdf_tdigest_name = "test_cdf_digest" + 
std::to_string(util::GetTimeStampMS());
+  bool exists = false;
+  auto status = tdigest_->Create(*ctx_, cdf_tdigest_name, {100}, &exists);
+  ASSERT_FALSE(exists);
+  ASSERT_TRUE(status.ok());
+
+  std::vector<double> samples = {1, 2, 2, 3, 3, 3, 4, 4, 4, 4, 5, 5, 5, 5, 5};
+  status = tdigest_->Add(*ctx_, cdf_tdigest_name, samples);
+  ASSERT_TRUE(status.ok());
+
+  std::vector<double> cdf_vals = {0, 1, 2, 3, 4, 5, 6};
+  redis::TDigestCDFResult result;
+
+  status = tdigest_->CDF(*ctx_, cdf_tdigest_name, cdf_vals, &result);
+  ASSERT_TRUE(status.ok()) << status.ToString();
+
+  std::vector<double> expected = {0.00, 0.03, 0.13, 0.29, 0.53, 0.83, 1.00};
+  ASSERT_TRUE(result.cdf_values) << "CDF should have values";
+  ASSERT_EQ(result.cdf_values->size(), cdf_vals.size());
+
+  for (size_t i = 0; i < cdf_vals.size(); i++) {
+    EXPECT_NEAR((*result.cdf_values)[i], expected[i], 0.015) << 
fmt::format("Mismatch at index {}", i);
+  }
+}
+
+TEST_F(RedisTDigestTest, CDF_returns_nan_on_empty_tdigest) {
+  std::string test_digest_name = "test_digest_cdf_nan" + 
std::to_string(util::GetTimeStampMS());
+
+  bool exists = false;
+  auto status = tdigest_->Create(*ctx_, test_digest_name, {100}, &exists);
+  ASSERT_FALSE(exists);
+  ASSERT_TRUE(status.ok());
+
+  std::vector<double> values = {0.0, 1.0, 2.0, 3.0};
+  redis::TDigestCDFResult result;
+
+  status = tdigest_->CDF(*ctx_, test_digest_name, values, &result);
+  ASSERT_TRUE(status.ok()) << status.ToString();
+  ASSERT_TRUE(result.cdf_values);
+}

Review Comment:
   `CDF_returns_nan_on_empty_tdigest` doesn't assert the externally observable 
behavior implied by its name (that each returned CDF value is NaN, and that the 
output length matches the number of inputs). This makes the test pass even if 
the implementation returns the wrong values.



##########
src/commands/cmd_tdigest.cc:
##########
@@ -556,6 +556,57 @@ class CommandTDigestTrimmedMean : public Commander {
   double high_cut_quantile_;
 };
 
+class CommandTDigestCDF : public Commander {
+  Status Parse(const std::vector<std::string> &args) override {
+    if (args.size() == 2) return {Status::RedisParseErr, 
errWrongNumOfArguments};
+    key_name_ = args[1];
+    values_.reserve(args.size() - 2);
+    for (size_t i = 2; i < args.size(); i++) {
+      auto value = ParseFloat(args[i]);
+      if (!value) {
+        return {Status::RedisParseErr, errValueIsNotFloat};
+      }
+      values_.push_back(*value);
+    }
+    return Status::OK();
+  }
+
+  Status Execute(engine::Context &ctx, Server *srv, Connection *conn, 
std::string *output) override {
+    TDigest tdigest(srv->storage, conn->GetNamespace());
+    std::vector<std::string> cdf_result;
+    TDigestCDFResult result;
+    TDigestMetadata metadata;
+    auto meta_status = tdigest.GetMetaData(ctx, key_name_, &metadata);
+    std::vector<std::string> nan_results(values_.size(), "nan");
+    if (!meta_status.ok()) {
+      if (meta_status.IsNotFound()) {
+        return {Status::RedisExecErr, errKeyNotFound};
+      }
+      *output = redis::MultiBulkString(RESP::v2, nan_results);
+      return Status::OK();
+    }
+    if (metadata.total_observations == 0) {
+      *output = redis::MultiBulkString(RESP::v2, nan_results);
+      return Status::OK();
+    }
+    auto s = tdigest.CDF(ctx, key_name_, values_, &result);
+    if (!s.ok()) {
+      return {Status::RedisExecErr, s.ToString()};
+    }
+    if (result.cdf_values) {
+      for (const auto &val : *result.cdf_values) {
+        cdf_result.push_back(util::Float2String(val));
+      }
+    }
+    *output = redis::MultiBulkString(RESP::v2, cdf_result);
+    return Status::OK();

Review Comment:
   The success reply is built with `redis::MultiBulkString(RESP::v2, ...)`, 
which bypasses `Connection` helpers used elsewhere in this file and may not 
honor the negotiated RESP mode.



##########
src/types/redis_tdigest.cc:
##########
@@ -570,6 +570,76 @@ rocksdb::Status TDigest::Merge(engine::Context& ctx, const 
Slice& dest_digest,
   return storage_->Write(ctx, storage_->DefaultWriteOptions(), 
batch->GetWriteBatch());
 }
 
+rocksdb::Status TDigest::CDF(engine::Context& ctx, const Slice& digest_name, 
const std::vector<double>& inputs,
+                             TDigestCDFResult* result) {
+  auto ns_key = AppendNamespacePrefix(digest_name);
+  TDigestMetadata metadata;
+  {
+    LockGuard guard(storage_->GetLockManager(), ns_key);
+
+    if (auto status = getMetaDataByNsKey(ctx, ns_key, &metadata); 
!status.ok()) {
+      return status;
+    }
+
+    if (metadata.unmerged_nodes > 0) {
+      auto batch = storage_->GetWriteBatchBase();
+      WriteBatchLogData log_data(kRedisTDigest);
+      if (auto status = batch->PutLogData(log_data.Encode()); !status.ok()) {
+        return status;
+      }
+
+      if (auto status = mergeCurrentBuffer(ctx, ns_key, batch, &metadata); 
!status.ok()) {
+        return status;
+      }
+      if (metadata.total_observations == 0) {
+        return rocksdb::Status::OK();
+      }

Review Comment:
   `TDigest::CDF` returns early when a buffer merge results in 
`total_observations == 0`, but it doesn't populate `result->cdf_values`. Other 
TDigest APIs consistently return a value vector (or rely on the optional being 
unset) to let callers produce per-input `nan` results for empty digests.



##########
src/commands/cmd_tdigest.cc:
##########
@@ -556,6 +556,57 @@ class CommandTDigestTrimmedMean : public Commander {
   double high_cut_quantile_;
 };
 
+class CommandTDigestCDF : public Commander {
+  Status Parse(const std::vector<std::string> &args) override {
+    if (args.size() == 2) return {Status::RedisParseErr, 
errWrongNumOfArguments};
+    key_name_ = args[1];
+    values_.reserve(args.size() - 2);
+    for (size_t i = 2; i < args.size(); i++) {
+      auto value = ParseFloat(args[i]);
+      if (!value) {
+        return {Status::RedisParseErr, errValueIsNotFloat};
+      }
+      values_.push_back(*value);
+    }
+    return Status::OK();
+  }
+
+  Status Execute(engine::Context &ctx, Server *srv, Connection *conn, 
std::string *output) override {
+    TDigest tdigest(srv->storage, conn->GetNamespace());
+    std::vector<std::string> cdf_result;
+    TDigestCDFResult result;
+    TDigestMetadata metadata;
+    auto meta_status = tdigest.GetMetaData(ctx, key_name_, &metadata);
+    std::vector<std::string> nan_results(values_.size(), "nan");
+    if (!meta_status.ok()) {
+      if (meta_status.IsNotFound()) {
+        return {Status::RedisExecErr, errKeyNotFound};
+      }
+      *output = redis::MultiBulkString(RESP::v2, nan_results);
+      return Status::OK();
+    }
+    if (metadata.total_observations == 0) {
+      *output = redis::MultiBulkString(RESP::v2, nan_results);
+      return Status::OK();
+    }

Review Comment:
   `TDIGEST.CDF` builds array replies using `redis::MultiBulkString(RESP::v2, 
...)`, which hard-codes RESP2 and is inconsistent with nearby TDigest commands 
that use `conn->MultiBulkString(...)` to respect the negotiated protocol.



##########
tests/gocase/unit/type/tdigest/tdigest_test.go:
##########
@@ -1309,4 +1310,145 @@ func tdigestTestsByRankAndByRevRank(t *testing.T, 
configs util.KvrocksServerConf
                        }
                }
        })
+
+       t.Run("tdigest.cdf with different arguments", func(t *testing.T) {
+               keyPrefix := "tdigest_cdf_"
+
+               require.ErrorContains(t, rdb.Do(ctx, "TDIGEST.CDF").Err(), 
errMsgWrongNumberArg)
+               require.ErrorContains(t, rdb.Do(ctx, "TDIGEST.CDF", 
keyPrefix+"key1").Err(), errMsgWrongNumberArg)
+
+               // non-existent key
+               require.ErrorContains(t, rdb.Do(ctx, "TDIGEST.CDF", 
keyPrefix+"nonexistent", "1.0").Err(), errMsgKeyNotExist)
+
+               // invalid float value
+               require.ErrorContains(t, rdb.Do(ctx, "TDIGEST.CDF", 
keyPrefix+"key2", "invalid").Err(), errValueIsNotFloat)
+
+               // create a tdigest and add some data
+               tdigestKey := keyPrefix + "source"
+               require.NoError(t, rdb.Do(ctx, "TDIGEST.CREATE", 
tdigestKey).Err())
+               require.NoError(t, rdb.Do(ctx, "TDIGEST.ADD", tdigestKey, 
"1.0", "2.0", "3.0", "4.0", "5.0").Err())
+
+               // single-value CDF query
+               rsp := rdb.Do(ctx, "TDIGEST.CDF", tdigestKey, "3.0")
+               require.NoError(t, rsp.Err())
+               vals, err := rsp.Slice()
+               require.NoError(t, err)
+               require.Len(t, vals, 1)
+               require.NotEqual(t, "nan", vals[0])
+
+               // multi-value CDF query
+               rsp = rdb.Do(ctx, "TDIGEST.CDF", tdigestKey, "0.0", "2.5", 
"5.0", "10.0")
+               require.NoError(t, rsp.Err())
+               vals, err = rsp.Slice()
+               require.NoError(t, err)
+               require.Len(t, vals, 4)
+
+               // empty tdigest should return "nan"
+               emptyKey := keyPrefix + "empty"
+               require.NoError(t, rdb.Do(ctx, "TDIGEST.CREATE", 
emptyKey).Err())
+               rsp = rdb.Do(ctx, "TDIGEST.CDF", emptyKey, "1.0")
+               require.NoError(t, rsp.Err())
+               vals, err = rsp.Slice()
+               require.NoError(t, err)
+               require.Len(t, vals, 1)
+               require.Equal(t, "nan", vals[0])
+
+               // testing with a empry digest with multi-valued CDF

Review Comment:
   Typo in comment: "empry" → "empty" (also adjust article to keep the sentence 
grammatical).



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