This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12013-b03cf1f1697ed84d0f57f416673cec40a075ef5a in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 1d18735bd38af023a33ae5ff77a738fcac55a549 Author: Goutam Adwant <[email protected]> AuthorDate: Mon Aug 31 09:17:53 2026 +0000 [Fix][Zeta] Stabilize pending job scheduler regression test (#12013) Signed-off-by: goutamadwant <[email protected]> --- .../seatunnel/engine/server/CoordinatorServiceTest.java | 11 +++++------ 1 file changed, 5 insertions(+), 6 deletions(-) diff --git a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java index 7e385e38f7..da1b2ae928 100644 --- a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java +++ b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/CoordinatorServiceTest.java @@ -623,6 +623,7 @@ public class CoordinatorServiceTest { EngineConfig engineConfig = new EngineConfig(); engineConfig.setScheduleStrategy(ScheduleStrategy.REJECT); CoordinatorService coordinatorService = newMockCoordinatorService(server, engineConfig); + CountDownLatch firstScheduleStarted = new CountDownLatch(1); CountDownLatch allowFirstScheduleToFinish = new CountDownLatch(1); try { JobMaster blockedJobMaster = @@ -630,6 +631,7 @@ public class CoordinatorServiceTest { Mockito.when(blockedJobMaster.preApplyResources()) .thenAnswer( invocation -> { + firstScheduleStarted.countDown(); allowFirstScheduleToFinish.await(); return true; }); @@ -637,12 +639,9 @@ public class CoordinatorServiceTest { ReflectionUtils.setField(coordinatorService, "isActive", true); invokePendingJobScheduler(coordinatorService); - // Wait until the scheduler thread is parked inside preApplyResources(). - await().atMost(5, TimeUnit.SECONDS) - .untilAsserted( - () -> - Mockito.verify(blockedJobMaster, Mockito.atLeastOnce()) - .preApplyResources()); + Assertions.assertTrue( + firstScheduleStarted.await(30, TimeUnit.SECONDS), + "Pending job scheduler did not enter preApplyResources"); // Simulate a master step-down. The blocked JobMaster is interrupted; the // PendingJobInfo must be dropped from the queue so a later restore cannot
