From 9de0b64c9ec229a932596d3e6081f200e2f90ddb Mon Sep 17 00:00:00 2001 From: Jark Wu Date: Thu, 20 Aug 2026 11:05:53 +0800 Subject: [PATCH] [test] Wait for local leader readiness Ensure waitAndGetLeader does not return until the assigned replica has completed its local leader transition. Co-Authored-By: Codex AI-Model: gpt-5.6-sol AI-Contributed/Feature: 0/0 AI-Contributed/UT: 15/15 --- .../server/testutils/FlussClusterExtension.java | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/fluss-server/src/test/java/org/apache/fluss/server/testutils/FlussClusterExtension.java b/fluss-server/src/test/java/org/apache/fluss/server/testutils/FlussClusterExtension.java index bc833e9c02b..0cc615482ab 100644 --- a/fluss-server/src/test/java/org/apache/fluss/server/testutils/FlussClusterExtension.java +++ b/fluss-server/src/test/java/org/apache/fluss/server/testutils/FlussClusterExtension.java @@ -943,8 +943,21 @@ public void waitUntilPartitionsDropped(TablePath tablePath, List dropped "Fail to wait partitions dropped"); } + /** Wait until the assigned replica is ready as the local leader and return its server id. */ public int waitAndGetLeader(TableBucket tb) { - return waitLeaderAndIsrReady(tb).leader(); + ZooKeeperClient zkClient = getZooKeeperClient(); + return waitValue( + () -> { + Optional leaderAndIsrOpt = zkClient.getLeaderAndIsr(tb); + if (!leaderAndIsrOpt.isPresent()) { + return Optional.empty(); + } + + int leader = leaderAndIsrOpt.get().leader(); + return getReplica(tb, leader, true).map(ignored -> leader); + }, + Duration.ofMinutes(1), + "Fail to wait leader replica ready for " + tb); } public int waitAndGetStandby(TableBucket tb) {