Skip to content

Commit b9d7dd3

Browse files
authored
[test] stabilize five confirmed flaky tests (#283)
* [meta] stabilize redis client pool expansion test * [meta] align reclaim sampling test with backend contract * [config] wait for follower leader discovery * [client] reserve deterministic bad grpc address * [integration_test] wait for reclaim metadata visibility
1 parent 8baaa45 commit b9d7dd3

6 files changed

Lines changed: 130 additions & 42 deletions

File tree

integration_test/reclaimer/reclaiming_test.py

Lines changed: 22 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -590,10 +590,11 @@ def _assert_no_over_eviction_with_delay(
590590
0,
591591
"the end-to-end Future must remain pending during the delay",
592592
)
593-
self.assertEqual(
594-
self._count_surviving_blocks(self._instance_id, range(12)),
595-
8,
596-
"CAS should hide exactly one batch while its Future is pending",
593+
self._wait_surviving_block_count(
594+
self._instance_id,
595+
range(12),
596+
expected_count=8,
597+
timeout_s=2,
597598
)
598599

599600
self._wait_metric_value(
@@ -724,6 +725,23 @@ def _count_surviving_blocks(self, instance_id, block_keys):
724725
surviving_blocks += 1
725726
return surviving_blocks
726727

728+
def _wait_surviving_block_count(
729+
self, instance_id, block_keys, expected_count, timeout_s
730+
):
731+
deadline = time.monotonic() + timeout_s
732+
last_count = None
733+
while time.monotonic() < deadline:
734+
last_count = self._count_surviving_blocks(
735+
instance_id, block_keys
736+
)
737+
if last_count == expected_count:
738+
return
739+
time.sleep(0.05)
740+
self.fail(
741+
f"surviving block count did not become {expected_count}; "
742+
f"last count: {last_count}"
743+
)
744+
727745
def test_persist_recover_00(self):
728746
"""Test e2e persist/recover: cache locations and metadata
729747
survive a normal server restart.

kv_cache_manager/client/src/internal/stub/test/grpc_stub_test.cc

Lines changed: 36 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
#include <string>
66
#include <sys/socket.h>
77
#include <thread>
8+
#include <unistd.h>
89

910
#include "kv_cache_manager/client/src/internal/stub/grpc_stub.h"
1011
#include "kv_cache_manager/common/logger.h"
@@ -44,6 +45,24 @@ bool WaitUntil(Func condition, int timeout_ms = 5000, int interval_ms = 100) {
4445
return false;
4546
}
4647

48+
class ScopedSocket {
49+
public:
50+
explicit ScopedSocket(int fd) : fd_(fd) {}
51+
~ScopedSocket() {
52+
if (fd_ >= 0) {
53+
close(fd_);
54+
}
55+
}
56+
57+
ScopedSocket(const ScopedSocket &) = delete;
58+
ScopedSocket &operator=(const ScopedSocket &) = delete;
59+
60+
int Get() const { return fd_; }
61+
62+
private:
63+
int fd_;
64+
};
65+
4766
} // namespace
4867

4968
class GrpcStubTest : public TESTBASE {
@@ -177,8 +196,24 @@ Stub::LocationSpecInfoMap GrpcStubTest::createLocationSpecInfos(int32_t spec_siz
177196
}
178197

179198
TEST_F(GrpcStubTest, TestBadAddress) {
199+
// Keep an ephemeral loopback port reserved without listening on it. This
200+
// guarantees connection refusal without racing another process for port_ + 1.
201+
ScopedSocket unavailable_socket(socket(AF_INET, SOCK_STREAM, 0));
202+
ASSERT_NE(-1, unavailable_socket.Get());
203+
204+
struct sockaddr_in addr;
205+
std::memset(&addr, 0, sizeof(addr));
206+
addr.sin_family = AF_INET;
207+
addr.sin_addr.s_addr = htonl(INADDR_LOOPBACK);
208+
addr.sin_port = 0;
209+
ASSERT_EQ(0, bind(unavailable_socket.Get(), reinterpret_cast<struct sockaddr *>(&addr), sizeof(addr)));
210+
211+
socklen_t len = sizeof(addr);
212+
ASSERT_EQ(0, getsockname(unavailable_socket.Get(), reinterpret_cast<struct sockaddr *>(&addr), &len));
213+
const int unavailable_port = ntohs(addr.sin_port);
214+
180215
stub_ = std::make_shared<GrpcStub>();
181-
ASSERT_EQ(ER_CONNECT_FAIL, stub_->AddConnection("0.0.0.0:" + std::to_string(port_ + 1), 1000));
216+
ASSERT_EQ(ER_CONNECT_FAIL, stub_->AddConnection("127.0.0.1:" + std::to_string(unavailable_port), 1000));
182217
}
183218

184219
TEST_F(GrpcStubTest, TestRetry) {

kv_cache_manager/config/test/leader_elector_test.cc

Lines changed: 19 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -654,21 +654,32 @@ TEST_F(LeaderElectorTest, LeaderDiscoveryTwoNodes) {
654654
std::this_thread::sleep_for(std::chrono::milliseconds(100));
655655
}
656656

657-
ASSERT_TRUE(elector1->IsLeader() || elector2->IsLeader());
657+
const bool is_elector1_leader = elector1->IsLeader();
658+
const bool is_elector2_leader = elector2->IsLeader();
659+
ASSERT_NE(is_elector1_leader, is_elector2_leader);
658660

659-
// 从任一 elector 获取 leader node id
660-
std::string leader_id = elector1->GetLeaderNodeID();
661-
EXPECT_FALSE(leader_id.empty());
661+
LeaderElector *leader = is_elector1_leader ? elector1.get() : elector2.get();
662+
LeaderElector *follower = is_elector1_leader ? elector2.get() : elector1.get();
663+
const std::string leader_id = leader->GetSelfNodeID();
662664

663-
// 用非 leader 的 elector 读取 leader 的节点信息(模拟从 follower 查询 leader)
664-
LeaderElector *follower = elector1->IsLeader() ? elector2.get() : elector1.get();
665+
// Election and follower discovery are separate work-loop observations. Wait
666+
// until the follower has observed the elected node before querying its info.
665667
NodeEndpointInfo leader_info;
666-
ErrorCode ec = follower->GetNodeInfo(leader_id, leader_info);
668+
ErrorCode ec = EC_NOENT;
669+
for (int i = 0; i < 100; ++i) {
670+
if (follower->GetLeaderNodeID() == leader_id) {
671+
ec = follower->GetLeaderNodeInfo(leader_info);
672+
if (ec == EC_OK) {
673+
break;
674+
}
675+
}
676+
std::this_thread::sleep_for(std::chrono::milliseconds(10));
677+
}
667678
ASSERT_EQ(EC_OK, ec);
668679

669680
// 验证读取到的信息和 leader 的注册信息一致
670681
EXPECT_EQ(leader_id, leader_info.node_id());
671-
if (elector1->IsLeader()) {
682+
if (is_elector1_leader) {
672683
EXPECT_EQ(host1, leader_info.host());
673684
EXPECT_EQ(8001, leader_info.meta_rpc_port());
674685
EXPECT_EQ(8002, leader_info.meta_http_port());

kv_cache_manager/meta/meta_storage_backend.h

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -418,7 +418,7 @@ class MetaStorageBackend {
418418
// 采样适合回收的 key(通常按 LRU 或 access time 排序)。
419419
// @param request_context 请求上下文;可为 nullptr
420420
// @param count 期望采样数量
421-
// @param out_keys [out] 采样结果
421+
// @param out_keys [out] 采样结果(实际数量可能小于 count)
422422
// @return EC_OK 成功;EC_ERROR 采样失败
423423
virtual ErrorCode
424424
SampleReclaimKeys(RequestContext *request_context, const int64_t count, KeyTypeVec &out_keys) noexcept = 0;

kv_cache_manager/meta/test/meta_indexer_test_base.cc

Lines changed: 14 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -252,21 +252,16 @@ void MetaIndexerTestBase::DoScanAndSampleReclaimKeysTest() {
252252
std::sort(keys.begin(), keys.end());
253253
ASSERT_EQ((KeyVector{0, 1, 2}), keys);
254254

255-
// 2. SampleReclaimKeys must converge to the full key set after enough tries.
256-
keys.clear();
257-
try_count = 100;
258-
while (try_count-- && keys.size() < static_cast<size_t>(key_count)) {
259-
KeyVector out_keys;
260-
ASSERT_EQ(EC_OK, meta_indexer_->SampleReclaimKeys(request_context_.get(), key_count, out_keys));
261-
for (const auto key : out_keys) {
262-
if (std::find(keys.begin(), keys.end(), key) == keys.end()) {
263-
keys.push_back(key);
264-
}
265-
}
255+
// 2. SampleReclaimKeys returns a non-empty subset. The requested count is
256+
// a hint: a sharded backend may return fewer keys when its per-shard
257+
// sampling quota is exhausted.
258+
KeyVector out_keys;
259+
ASSERT_EQ(EC_OK, meta_indexer_->SampleReclaimKeys(request_context_.get(), key_count, out_keys));
260+
ASSERT_FALSE(out_keys.empty());
261+
ASSERT_LE(out_keys.size(), static_cast<size_t>(key_count));
262+
for (const auto key : out_keys) {
263+
ASSERT_TRUE(key >= 0 && key < key_count) << "Unexpected key: " << key;
266264
}
267-
ASSERT_GT(try_count, 0);
268-
std::sort(keys.begin(), keys.end());
269-
ASSERT_EQ((KeyVector{0, 1, 2}), keys);
270265

271266
// 3. Cleanup
272267
meta_indexer_->Delete(request_context_.get(), data.keys);
@@ -453,11 +448,11 @@ void MetaIndexerTestBase::DoReadModifyWriteLocationTest() {
453448
}
454449

455450
void MetaIndexerTestBase::DoSimpleTest() {
456-
DoPutTest();
457-
DoDeleteAndExistTest();
458-
DoScanAndSampleReclaimKeysTest();
459-
DoReadModifyWriteBlockTest();
460-
DoReadModifyWriteLocationTest();
451+
ASSERT_NO_FATAL_FAILURE(DoPutTest());
452+
ASSERT_NO_FATAL_FAILURE(DoDeleteAndExistTest());
453+
ASSERT_NO_FATAL_FAILURE(DoScanAndSampleReclaimKeysTest());
454+
ASSERT_NO_FATAL_FAILURE(DoReadModifyWriteBlockTest());
455+
ASSERT_NO_FATAL_FAILURE(DoReadModifyWriteLocationTest());
461456
}
462457

463458
void MetaIndexerTestBase::DoMultiThreadTest() {

kv_cache_manager/meta/test/meta_redis_backend_test.cc

Lines changed: 38 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,3 +1,6 @@
1+
#include <atomic>
2+
#include <future>
3+
#include <mutex>
14
#include <thread>
25

36
#include "kv_cache_manager/common/redis_client.h"
@@ -319,7 +322,20 @@ TEST_F(MetaRedisBackendTest, TestRedisError) {
319322
}
320323

321324
TEST_F(MetaRedisBackendTest, TestMultiThreadSimple) {
322-
auto make_mock_client = []() {
325+
std::atomic<int> client_index{0};
326+
std::once_flag first_pipeline_call;
327+
std::promise<void> first_client_entered_promise;
328+
auto first_client_entered = first_client_entered_promise.get_future();
329+
std::promise<void> release_first_client_promise;
330+
auto release_first_client = release_first_client_promise.get_future().share();
331+
std::promise<void> second_client_created_promise;
332+
auto second_client_created = second_client_created_promise.get_future();
333+
334+
auto make_mock_client = [&]() {
335+
const int current_client_index = client_index.fetch_add(1);
336+
if (current_client_index == 1) {
337+
second_client_created_promise.set_value();
338+
}
323339
StandardUri empty_storage_uri;
324340
auto mock_redis_client = std::make_unique<MockRedisClient>(empty_storage_uri);
325341
EXPECT_CALL(*mock_redis_client, Reconnect()).WillRepeatedly(Return(true));
@@ -329,8 +345,11 @@ TEST_F(MetaRedisBackendTest, TestMultiThreadSimple) {
329345
TryExecPipeline(ElementsAre(
330346
ElementsAre(StrEq("HMGET"), StrEq("kvcache:instance_instance_0:cache_1"), StrEq("f1"), StrEq("f2")),
331347
ElementsAre(StrEq("HMGET"), StrEq("kvcache:instance_instance_0:cache_2"), StrEq("f1"), StrEq("f2")))))
332-
.WillRepeatedly(Invoke([]() {
333-
usleep(5 * 1000); // assume network use 5ms
348+
.WillRepeatedly(Invoke([&, current_client_index]() {
349+
if (current_client_index == 0) {
350+
std::call_once(first_pipeline_call, [&] { first_client_entered_promise.set_value(); });
351+
release_first_client.wait();
352+
}
334353
std::vector<ReplyUPtr> get_replies_2;
335354
get_replies_2.emplace_back(MakeFakeReplyArrayString({"v1-1", "v1-2"}));
336355
get_replies_2.emplace_back(MakeFakeReplyArrayString({"v2-1", "v2-2"}));
@@ -353,14 +372,24 @@ TEST_F(MetaRedisBackendTest, TestMultiThreadSimple) {
353372
{EC_OK, EC_OK},
354373
{{{"f1", "v1-1"}, {"f2", "v1-2"}}, {{"f1", "v2-1"}, {"f2", "v2-2"}}});
355374
};
356-
std::vector<std::thread> threads;
357-
for (int i = 0; i < 10; ++i) {
358-
threads.emplace_back(get_task);
359-
usleep(2 * 1000);
375+
std::thread first_thread(get_task);
376+
const auto first_entered_status = first_client_entered.wait_for(std::chrono::seconds(1));
377+
378+
std::thread second_thread;
379+
std::future_status second_created_status = std::future_status::timeout;
380+
if (first_entered_status == std::future_status::ready) {
381+
second_thread = std::thread(get_task);
382+
second_created_status = second_client_created.wait_for(std::chrono::seconds(1));
360383
}
361-
for (auto &thread : threads) {
362-
thread.join();
384+
385+
release_first_client_promise.set_value();
386+
first_thread.join();
387+
if (second_thread.joinable()) {
388+
second_thread.join();
363389
}
390+
391+
ASSERT_EQ(std::future_status::ready, first_entered_status);
392+
ASSERT_EQ(std::future_status::ready, second_created_status);
364393
ASSERT_EQ(EC_OK, meta_redis_backend_->Close());
365394
}
366395

0 commit comments

Comments
 (0)