This is an automated email from the ASF dual-hosted git repository.

davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new f13b610eff [Fix][Zeta] Avoid duplicate pending job scheduling after 
failover (#11653)
f13b610eff is described below

commit f13b610effa5769c79e1903969ee4869525fb9a9
Author: Jast <[email protected]>
AuthorDate: Sun Aug 23 16:16:50 2026 +0800

    [Fix][Zeta] Avoid duplicate pending job scheduling after failover (#11653)
    
    Co-authored-by: David Zollo <[email protected]>
---
 .../engine/server/CoordinatorService.java          | 278 +++++++++++++++------
 .../engine/server/utils/PeekBlockingQueue.java     |  63 ++++-
 .../engine/server/CoordinatorServiceTest.java      | 237 ++++++++++++++++++
 3 files changed, 498 insertions(+), 80 deletions(-)

diff --git 
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java
 
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java
index 15c76c8835..328d245095 100644
--- 
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java
+++ 
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java
@@ -122,6 +122,7 @@ import java.util.concurrent.TimeUnit;
 import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.locks.ReentrantLock;
 import java.util.stream.Collectors;
 
 import static 
org.apache.seatunnel.api.options.EnvCommonOptions.CHECKPOINT_INTERVAL;
@@ -185,6 +186,25 @@ public class CoordinatorService {
     private final PeekBlockingQueue<PendingJobInfo> pendingJobQueue =
             new PeekBlockingQueue<>(PendingJobInfo::getJobId);
 
+    // Fair lock that serializes "peek the head job and reserve it for 
evaluation" across all
+    // scheduler generations. Holding this lock around peek+reserve guarantees 
that at most one
+    // thread holds a given job id in schedulingPendingJobIds at any time, so 
two overlapping
+    // scheduler threads (e.g. a stale one still inside preApplyResources and 
a fresh one started
+    // after a master flap) cannot independently reserve the same job.
+    private final ReentrantLock pendingJobScheduleLock = new 
ReentrantLock(true);
+    // Monotonically increasing generation counter for scheduler runs. Each 
scheduler thread
+    // snapshots this value at submission time and exits when it is bumped. 
Bumped inside
+    // checkNewActiveMaster and clearCoordinatorService to invalidate any 
in-flight scheduler.
+    private final AtomicLong pendingJobScheduleEpoch = new AtomicLong();
+    // JobMasters currently inside jobMaster.preApplyResources(). Interrupted 
en masse by
+    // clearCoordinatorService so they release any pending RPC and stop 
touching shared state.
+    private final Set<JobMaster> schedulingJobMasters = 
ConcurrentHashMap.newKeySet();
+    // Job ids currently reserved by a scheduler thread for evaluation under
+    // pendingJobScheduleLock. reservePendingJobInfo adds; 
releasePendingJobInfo removes. Cleared
+    // by clearCoordinatorService so a freshly-restored scheduler can 
re-reserve jobs after a
+    // master flap.
+    private final Set<Long> schedulingPendingJobIds = 
ConcurrentHashMap.newKeySet();
+
     /**
      * IMap key is {@link PipelineLocation}
      *
@@ -286,21 +306,27 @@ public class CoordinatorService {
      */
     private void startPendingJobScheduleThread() {
         logger.info("Start pending job schedule thread");
+        long scheduleEpoch = pendingJobScheduleEpoch.get();
         Runnable pendingJobScheduleTask =
                 () -> {
                     
Thread.currentThread().setName("pending-job-schedule-runner");
-                    while (isActive) {
+                    while (isPendingJobSchedulerCurrent(scheduleEpoch)) {
                         try {
-                            pendingJobSchedule();
-                        } catch (InterruptedException interrupted) {
-                            throw new RuntimeException(interrupted);
+                            pendingJobSchedule(scheduleEpoch);
+                        } catch (InterruptedException e) {
+                            Thread.currentThread().interrupt();
+                            return;
                         } catch (Throwable e) {
+                            if (!isPendingJobSchedulerCurrent(scheduleEpoch)) {
+                                return;
+                            }
                             logger.severe("Error in pending job schedule 
thread", e);
                             try {
                                 Thread.sleep(3000L);
                             } catch (InterruptedException ex) {
                                 logger.severe("Pending job schedule thread 
interrupted", ex);
                                 Thread.currentThread().interrupt();
+                                return;
                             }
                         }
                     }
@@ -308,6 +334,50 @@ public class CoordinatorService {
         executorService.submit(pendingJobScheduleTask);
     }
 
+    /**
+     * A scheduler generation is "current" only when this node is still the 
active master AND no
+     * newer generation has been started since this thread was submitted. 
Every scheduler-loop
+     * iteration and every stale-branch decision hinges on this predicate: 
once it turns false the
+     * thread must stop touching shared scheduling state, because a newer 
generation now owns it.
+     */
+    private boolean isPendingJobSchedulerCurrent(long scheduleEpoch) {
+        return isActive && pendingJobScheduleEpoch.get() == scheduleEpoch;
+    }
+
+    /**
+     * Reserves the next pending job for evaluation under 
pendingJobScheduleLock. The returned
+     * PendingJobInfo has its job id added to schedulingPendingJobIds, 
blocking any other scheduler
+     * thread from reserving the same job until {@link #releasePendingJobInfo} 
is called. Returns
+     * null only if the queue is empty.
+     */
+    private PendingJobInfo reservePendingJobInfo() throws InterruptedException 
{
+        pendingJobScheduleLock.lockInterruptibly();
+        try {
+            PendingJobInfo pendingJobInfo =
+                    pendingJobQueue.peekBlocking(
+                            jobInfo -> 
!schedulingPendingJobIds.contains(jobInfo.getJobId()));
+            if (Objects.isNull(pendingJobInfo)) {
+                return null;
+            }
+            schedulingPendingJobIds.add(pendingJobInfo.getJobId());
+            return pendingJobInfo;
+        } finally {
+            pendingJobScheduleLock.unlock();
+        }
+    }
+
+    /**
+     * Releases a previously reserved pending job and wakes any other 
scheduler thread that may be
+     * blocked inside reservePendingJobInfo waiting for a non-empty queue. 
Idempotent: safe to call
+     * when the job was never reserved.
+     */
+    private void releasePendingJobInfo(PendingJobInfo pendingJobInfo) {
+        if (Objects.nonNull(pendingJobInfo)
+                && schedulingPendingJobIds.remove(pendingJobInfo.getJobId())) {
+            pendingJobQueue.release();
+        }
+    }
+
     /**
      * Attempts to advance the head job in the pending queue into execution.
      *
@@ -317,10 +387,19 @@ public class CoordinatorService {
      * {@link JobMaster#run()}. Restored jobs rehydrate pipeline execution 
state before the new
      * master takes over execution.
      *
+     * <p>If {@code scheduleEpoch} no longer matches {@link 
#pendingJobScheduleEpoch} when this
+     * method finishes evaluating resources (because the node stepped down or 
master ownership
+     * transferred), the job's {@link JobMaster} is interrupted, its entry is 
removed from the
+     * pending queue, and the method returns without invoking {@code run()}. A 
new scheduler thread
+     * started for the current epoch will pick up (or rebuild) the job from 
{@code
+     * runningJobInfoIMap} instead of re-running the poisoned instance.
+     *
      * @throws InterruptedException if queue waiting or retry sleep is 
interrupted
+     * @throws Throwable any unchecked failure from {@code preApplyResources}; 
the reservation is
+     *     released before propagation so a future scheduler iteration can 
retry.
      */
-    private void pendingJobSchedule() throws InterruptedException {
-        PendingJobInfo pendingJobInfo = pendingJobQueue.peekBlocking();
+    private void pendingJobSchedule(long scheduleEpoch) throws Throwable {
+        PendingJobInfo pendingJobInfo = reservePendingJobInfo();
         if (Objects.isNull(pendingJobInfo)) {
             // This situation almost never happens because pendingJobSchedule 
is single-threaded
             logger.warning("The peek job info is null");
@@ -337,71 +416,88 @@ public class CoordinatorService {
                 String.format(
                         "Start calculating whether pending task resources are 
enough: %s", jobId));
 
-        boolean preApplyResources = jobMaster.preApplyResources();
-        if (!preApplyResources) {
-            try {
-                PendingJobDiagnostic diagnostic =
-                        PendingDiagnosticsCollector.collectJobDiagnostic(
-                                pendingJobInfo, Collections.emptyMap(), 
getResourceManager());
-                pendingJobInfo.recordSnapshot(diagnostic);
-            } catch (Exception e) {
-                logger.warning(
-                        String.format(
-                                "Collect pending diagnostic for job %s failed: 
%s",
-                                jobId, ExceptionUtils.getMessage(e)));
+        boolean preApplyResources;
+        schedulingJobMasters.add(jobMaster);
+        try {
+            preApplyResources = jobMaster.preApplyResources();
+        } catch (Throwable e) {
+            releasePendingJobInfo(pendingJobInfo);
+            throw e;
+        } finally {
+            schedulingJobMasters.remove(jobMaster);
+        }
+        try {
+            if (!isPendingJobSchedulerCurrent(scheduleEpoch)) {
+                jobMaster.interrupt();
+                // Drop the poisoned PendingJobInfo from the queue: its 
JobMaster has been
+                // interrupted (which permanently completes 
jobMasterCompleteFuture exceptionally),
+                // so the next-epoch scheduler must not re-dispatch it. 
Restore on the next
+                // activation will rebuild a fresh PendingJobInfo from 
runningJobInfoIMap.
+                pendingJobQueue.remove(pendingJobInfo);
+                return;
             }
-            logger.info(
-                    String.format(
-                            "Current strategy is %s, and resources is not 
enough, skipping this schedule, JobID: %s",
-                            scheduleStrategy, jobId));
-            if (isWaitStrategy) {
+            if (!preApplyResources) {
                 try {
-                    Thread.sleep(3000);
-                } catch (InterruptedException e) {
-                    logger.severe(ExceptionUtils.getMessage(e));
+                    PendingJobDiagnostic diagnostic =
+                            PendingDiagnosticsCollector.collectJobDiagnostic(
+                                    pendingJobInfo, Collections.emptyMap(), 
getResourceManager());
+                    pendingJobInfo.recordSnapshot(diagnostic);
+                } catch (Exception e) {
+                    logger.warning(
+                            String.format(
+                                    "Collect pending diagnostic for job %s 
failed: %s",
+                                    jobId, ExceptionUtils.getMessage(e)));
                 }
-                return;
-            } else {
-                completeFailJob(jobMaster);
-                queueRemove(jobMaster);
-                return;
-            }
-        }
-        logger.info(String.format("Resources enough, start running: %s", 
jobId));
-        // When deleting jobmaster from pendingJobQueue, make sure that there 
is a corresponding
-        // jobMaster in the runningJobMasterMap
-        runningJobMasterMap.put(jobId, jobMaster);
-        final PendingJobInfo finalPendingJobInfo = pendingJobQueue.take();
-        final JobMaster finalJobMaster = finalPendingJobInfo.getJobMaster();
-        PendingSourceState pendingSourceState = 
finalPendingJobInfo.getPendingSourceState();
-        MDCExecutorService mdcExecutorService = MDCTracer.tracing(jobId, 
executorService);
-        mdcExecutorService.submit(
-                () -> {
+                logger.info(
+                        String.format(
+                                "Current strategy is %s, and resources is not 
enough, skipping this schedule, JobID: %s",
+                                scheduleStrategy, jobId));
+                if (isWaitStrategy) {
                     try {
-                        String jobFullName = 
finalJobMaster.getPhysicalPlan().getJobFullName();
-                        JobStatus jobStatus = (JobStatus) 
runningJobStateIMap.get(jobId);
-                        if (pendingSourceState == PendingSourceState.RESTORE
-                                && !jobStatus.isEndState()) {
-                            finalJobMaster
-                                    .getPhysicalPlan()
-                                    .getPipelineList()
-                                    .forEach(SubPlan::restorePipelineState);
-                        }
-                        logger.info(
-                                String.format(
-                                        "The %s %s is in %s state, restore 
pipeline and take over this job running",
-                                        pendingSourceState, jobFullName, 
jobStatus));
-                        finalJobMaster.run();
-                    } finally {
-                        if (jobMasterCompletedSuccessfully(finalJobMaster, 
pendingSourceState)) {
-                            runningJobMasterMap.remove(jobId);
-                        }
+                        Thread.sleep(3000);
+                    } catch (InterruptedException e) {
+                        logger.severe(ExceptionUtils.getMessage(e));
                     }
-                });
-    }
-
-    private void queueRemove(JobMaster jobMaster) {
-        pendingJobQueue.removeById(jobMaster.getJobId());
+                    return;
+                } else {
+                    completeFailJob(jobMaster);
+                    pendingJobQueue.remove(pendingJobInfo);
+                    return;
+                }
+            }
+            logger.info(String.format("Resources enough, start running: %s", 
jobId));
+            // When deleting jobmaster from pendingJobQueue, make sure that 
there is a corresponding
+            // jobMaster in the runningJobMasterMap
+            runningJobMasterMap.put(jobId, jobMaster);
+            pendingJobQueue.remove(pendingJobInfo);
+            PendingSourceState pendingSourceState = 
pendingJobInfo.getPendingSourceState();
+            MDCExecutorService mdcExecutorService = MDCTracer.tracing(jobId, 
executorService);
+            mdcExecutorService.submit(
+                    () -> {
+                        try {
+                            String jobFullName = 
jobMaster.getPhysicalPlan().getJobFullName();
+                            JobStatus jobStatus = (JobStatus) 
runningJobStateIMap.get(jobId);
+                            if (pendingSourceState == 
PendingSourceState.RESTORE
+                                    && !jobStatus.isEndState()) {
+                                jobMaster
+                                        .getPhysicalPlan()
+                                        .getPipelineList()
+                                        
.forEach(SubPlan::restorePipelineState);
+                            }
+                            logger.info(
+                                    String.format(
+                                            "The %s %s is in %s state, restore 
pipeline and take over this job running",
+                                            pendingSourceState, jobFullName, 
jobStatus));
+                            jobMaster.run();
+                        } finally {
+                            if (jobMasterCompletedSuccessfully(jobMaster, 
pendingSourceState)) {
+                                runningJobMasterMap.remove(jobId);
+                            }
+                        }
+                    });
+        } finally {
+            releasePendingJobInfo(pendingJobInfo);
+        }
     }
 
     /**
@@ -1129,8 +1225,20 @@ public class CoordinatorService {
      * <p>When this node becomes the active master, the coordinator 
initializes distributed services
      * and triggers job restore. When it loses master ownership, local 
coordinator state is torn
      * down. Initialization failures are cleaned up locally and retried by 
later polling cycles.
+     *
+     * <p>Synchronized on the same monitor as {@link 
#clearCoordinatorService()} so that the "check
+     * ownership then flip {@code isActive}" sequence is atomic. Two threads 
running this method
+     * concurrently could otherwise both observe {@code isActive == false} and 
both perform the
+     * activation block, bumping {@link #pendingJobScheduleEpoch} twice 
without any step-down in
+     * between. The scheduler thread started by the first activation would 
then see its own epoch as
+     * stale once {@code preApplyResources} returns and discard the pending 
job it had already
+     * reserved (interrupt its JobMaster and drop it from {@code 
pendingJobQueue}) even though this
+     * node never actually lost master ownership. Because no {@code 
clearCoordinatorService()} ran,
+     * {@code restoreAllRunningJobFromMasterNodeSwitch} is never triggered to 
rebuild that entry, so
+     * the job would be silently lost. Epoch changes must therefore only ever 
come from a real
+     * activation or a real step-down.
      */
-    private void checkNewActiveMaster() {
+    private synchronized void checkNewActiveMaster() {
         try {
             if (!isActive && this.seaTunnelServer.isMasterNode()) {
                 logger.info(
@@ -1139,6 +1247,7 @@ public class CoordinatorService {
                     this.executorService = createCoordinatorExecutor();
                 }
                 initCoordinatorService();
+                pendingJobScheduleEpoch.incrementAndGet();
                 isActive = true;
                 startPendingJobScheduleThread();
                 seaTunnelServer.startRealtimeMetricsService(this);
@@ -1174,19 +1283,32 @@ public class CoordinatorService {
         if (!coordinatorServiceCleared.compareAndSet(false, true)) {
             return;
         }
+        pendingJobScheduleEpoch.incrementAndGet();
+        schedulingJobMasters.forEach(JobMaster::interrupt);
+        schedulingJobMasters.clear();
+        schedulingPendingJobIds.clear();
+        pendingJobQueue.release();
         // interrupt all JobMaster
         runningJobMasterMap.values().forEach(JobMaster::interrupt);
-        if (isWaitStrategy) {
-            pendingJobQueue
-                    .getJobIdMap()
-                    .values()
-                    .forEach(
-                            pendingJobInfo -> {
-                                JobMaster jobMaster = 
pendingJobInfo.getJobMaster();
-                                jobMaster.interrupt();
-                            });
-            pendingJobQueue.clear();
-        }
+        // Interrupt and discard every JobMaster currently sitting in 
pendingJobQueue. This is
+        // intentionally unconditional (not gated on isWaitStrategy): every 
entry in
+        // pendingJobQueue corresponds to a JobInfo already persisted in 
runningJobInfoIMap, so
+        // clearCoordinatorService can safely drop the local queue copy. 
Restore on the next
+        // activation will rebuild PendingJobInfo entries from 
runningJobInfoIMap via
+        // restoreAllRunningJobFromMasterNodeSwitch, producing a single fresh 
JobMaster per job
+        // id. Leaving the stale PendingJobInfo in the queue (the pre-fix 
behavior for
+        // REJECT strategy) caused an interrupted JobMaster's 
permanently-poisoned
+        // jobMasterCompleteFuture to be picked by the next-epoch scheduler 
and led to duplicate
+        // physicalPlan.startJob() invocations for the same job id.
+        pendingJobQueue
+                .getJobIdMap()
+                .values()
+                .forEach(
+                        pendingJobInfo -> {
+                            JobMaster jobMaster = 
pendingJobInfo.getJobMaster();
+                            jobMaster.interrupt();
+                        });
+        pendingJobQueue.clear();
         executorService.shutdownNow();
         runningJobMasterMap.clear();
 
diff --git 
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/utils/PeekBlockingQueue.java
 
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/utils/PeekBlockingQueue.java
index 699f5bca08..b866d6fbdb 100644
--- 
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/utils/PeekBlockingQueue.java
+++ 
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/utils/PeekBlockingQueue.java
@@ -29,6 +29,7 @@ import java.util.concurrent.locks.Condition;
 import java.util.concurrent.locks.Lock;
 import java.util.concurrent.locks.ReentrantLock;
 import java.util.function.Function;
+import java.util.function.Predicate;
 
 /**
  * PeekBlockingQueue implements blocking when peeking. Queues like 
BlockingQueue only support
@@ -76,18 +77,57 @@ public class PeekBlockingQueue<E> {
         return element;
     }
 
+    /**
+     * Wakes up any thread currently blocked inside {@link 
#peekBlocking(Predicate)} waiting for the
+     * queue to become non-empty. No-op when the queue is empty; a thread 
waiting on an
+     * actually-empty queue still relies on the caller's {@code 
shutdownNow}/{@code interrupt} to
+     * unblock, so do not remove the {@code shutdownNow} call in {@code
+     * CoordinatorService.clearCoordinatorService} assuming this method is 
sufficient on its own.
+     */
+    public void release() {
+        lock.lock();
+        try {
+            if (!queue.isEmpty()) {
+                notEmpty.signalAll();
+            }
+        } finally {
+            lock.unlock();
+        }
+    }
+
     public E peekBlocking() throws InterruptedException {
+        return peekBlocking(element -> true);
+    }
+
+    /**
+     * Blocks until the queue contains an element satisfying {@code predicate} 
(in FIFO order) and
+     * returns it without removing. The caller is responsible for removing the 
element via {@link
+     * #remove(Object)} (or {@link #take()}) once the element is no longer 
needed so the next waiter
+     * can make progress on a different element.
+     */
+    public E peekBlocking(Predicate<E> predicate) throws InterruptedException {
         lock.lock();
         try {
-            while (queue.peek() == null) {
+            E element = findFirst(predicate);
+            while (element == null) {
                 notEmpty.await();
+                element = findFirst(predicate);
             }
-            return queue.peek();
+            return element;
         } finally {
             lock.unlock();
         }
     }
 
+    private E findFirst(Predicate<E> predicate) {
+        for (E element : queue) {
+            if (predicate.test(element)) {
+                return element;
+            }
+        }
+        return null;
+    }
+
     public Integer size() {
         lock.lock();
         try {
@@ -124,6 +164,25 @@ public class PeekBlockingQueue<E> {
         }
     }
 
+    /**
+     * Removes a specific element by identity from both the underlying queue 
and the id-to-element
+     * map. Uses {@link Map#remove(Object, Object)} on the id map so a newer 
entry that happens to
+     * share the same id (e.g. a freshly restored {@code PendingJobInfo} that 
has overwritten the id
+     * pointer) is not accidentally evicted alongside the intended removal.
+     */
+    public boolean remove(E element) {
+        lock.lock();
+        try {
+            boolean removed = queue.remove(element);
+            if (removed) {
+                jobIdMap.remove(idExtractor.apply(element), element);
+            }
+            return removed;
+        } finally {
+            lock.unlock();
+        }
+    }
+
     public boolean contains(Long jobId) {
         return jobIdMap.containsKey(jobId);
     }
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 286ba67fa6..7e385e38f7 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
@@ -87,6 +87,7 @@ import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.CompletionException;
 import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.CyclicBarrier;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
@@ -95,6 +96,7 @@ import java.util.concurrent.ThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
 
 import static 
org.apache.seatunnel.engine.core.classloader.DefaultClassLoaderService.SKIP_CHECK_JAR;
 import static org.awaitility.Awaitility.await;
@@ -418,6 +420,79 @@ public class CoordinatorServiceTest {
         }
     }
 
+    /**
+     * Concurrent {@code checkNewActiveMaster()} calls must activate the 
coordinator exactly once.
+     *
+     * <p>Regression test for a race in which two threads both observed {@code 
isActive == false}
+     * and both ran the activation block, bumping the pending-job schedule 
epoch twice without any
+     * step-down in between. The scheduler thread started by the first 
activation then saw its own
+     * epoch as stale after {@code preApplyResources()} returned and discarded 
the job it had
+     * already reserved — interrupting its JobMaster and dropping it from 
{@code pendingJobQueue} —
+     * even though the node never lost master ownership. Because no {@code
+     * clearCoordinatorService()} ran, nothing rebuilt that entry, so the job 
was silently lost and
+     * never ran. This surfaced as a 30s {@code ConditionTimeoutException} in 
{@link
+     * #testFailoverStopsOldPendingQueueAndNewCoordinatorCanSchedule()}, whose 
reflective {@code
+     * checkNewActiveMaster()} call could interleave with the 
constructor-scheduled
+     * masterActiveListener tick.
+     */
+    @Test
+    void testConcurrentMasterActivationActivatesOnceAndKeepsPendingJob() 
throws Exception {
+        // Start as non-master so the constructor-scheduled 
masterActiveListener cannot activate
+        // before we stop it; that keeps this test's two threads the only 
activation drivers.
+        AtomicBoolean masterFlag = new AtomicBoolean(false);
+        SeaTunnelServer server = Mockito.mock(SeaTunnelServer.class);
+        Mockito.when(server.isMasterNode()).thenAnswer(invocation -> 
masterFlag.get());
+        CoordinatorService coordinatorService = 
newMockCoordinatorService(server);
+        int activationThreads = 4;
+        ExecutorService activationPool = 
Executors.newFixedThreadPool(activationThreads);
+        try {
+            CountDownLatch runLatch = new CountDownLatch(1);
+            JobMaster pendingJob = enqueueMockPendingJob(coordinatorService, 
80001L, runLatch);
+            long epochBefore = 
getPendingJobScheduleEpoch(coordinatorService).get();
+
+            masterFlag.set(true);
+            CyclicBarrier startTogether = new CyclicBarrier(activationThreads);
+            List<Future<?>> activations = new ArrayList<>();
+            for (int i = 0; i < activationThreads; i++) {
+                activations.add(
+                        activationPool.submit(
+                                () -> {
+                                    startTogether.await(30, TimeUnit.SECONDS);
+                                    
invokeCheckNewActiveMaster(coordinatorService);
+                                    return null;
+                                }));
+            }
+            for (Future<?> activation : activations) {
+                activation.get(60, TimeUnit.SECONDS);
+            }
+
+            await().atMost(30, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () -> 
Assertions.assertTrue(coordinatorService.isCoordinatorActive()));
+            Assertions.assertEquals(
+                    epochBefore + 1,
+                    getPendingJobScheduleEpoch(coordinatorService).get(),
+                    "concurrent checkNewActiveMaster calls must bump the 
schedule epoch only once");
+            await().atMost(30, TimeUnit.SECONDS)
+                    .untilAsserted(() -> Assertions.assertEquals(0L, 
runLatch.getCount()));
+            Mockito.verify(pendingJob, Mockito.times(1)).run();
+            Mockito.verify(pendingJob, Mockito.never()).interrupt();
+            
Assertions.assertFalse(coordinatorService.getPendingJobQueue().contains(80001L));
+        } finally {
+            activationPool.shutdownNow();
+            shutdownCoordinatorIfRunning(coordinatorService);
+        }
+    }
+
+    private AtomicLong getPendingJobScheduleEpoch(CoordinatorService 
coordinatorService) {
+        return ReflectionUtils.getField(coordinatorService, 
"pendingJobScheduleEpoch")
+                .map(AtomicLong.class::cast)
+                .orElseThrow(
+                        () ->
+                                new AssertionError(
+                                        "Failed to get pendingJobScheduleEpoch 
by reflection"));
+    }
+
     @Test
     void testCheckNewActiveMasterIsIdempotentWhenAlreadyActive() throws 
Exception {
         AtomicBoolean masterFlag = new AtomicBoolean(true);
@@ -444,6 +519,149 @@ public class CoordinatorServiceTest {
         }
     }
 
+    @Test
+    void testPendingJobSchedulerIgnoresJobReservedByPreviousScheduler() throws 
Exception {
+        SeaTunnelServer server = Mockito.mock(SeaTunnelServer.class);
+        CoordinatorService coordinatorService = 
newMockCoordinatorService(server);
+        CountDownLatch firstScheduleStarted = new CountDownLatch(1);
+        CountDownLatch allowFirstScheduleToFinish = new CountDownLatch(1);
+        try {
+            CountDownLatch runLatch = new CountDownLatch(1);
+            JobMaster jobMaster = enqueueMockPendingJob(coordinatorService, 
70001L, runLatch);
+            AtomicBoolean firstSchedule = new AtomicBoolean(true);
+            Mockito.when(jobMaster.preApplyResources())
+                    .thenAnswer(
+                            invocation -> {
+                                if (firstSchedule.compareAndSet(true, false)) {
+                                    firstScheduleStarted.countDown();
+                                    allowFirstScheduleToFinish.await();
+                                }
+                                return true;
+                            });
+
+            ReflectionUtils.setField(coordinatorService, "isActive", true);
+            invokePendingJobScheduler(coordinatorService);
+            Assertions.assertTrue(firstScheduleStarted.await(5, 
TimeUnit.SECONDS));
+
+            ReflectionUtils.setField(coordinatorService, "isActive", false);
+            ReflectionUtils.setField(coordinatorService, "isActive", true);
+            invokePendingJobScheduler(coordinatorService);
+
+            await().during(500, TimeUnit.MILLISECONDS)
+                    .atMost(1, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () -> Mockito.verify(jobMaster, 
Mockito.times(1)).preApplyResources());
+
+            allowFirstScheduleToFinish.countDown();
+            await().atMost(5, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () ->
+                                    Assertions.assertFalse(
+                                            coordinatorService
+                                                    .getPendingJobQueue()
+                                                    .contains(70001L)));
+            Mockito.verify(jobMaster, Mockito.atMostOnce()).run();
+        } finally {
+            allowFirstScheduleToFinish.countDown();
+            shutdownCoordinatorIfRunning(coordinatorService);
+        }
+    }
+
+    @Test
+    void 
testPendingJobSchedulerCanAdvanceNextJobWhenPreviousResourceCheckBlocks()
+            throws Exception {
+        SeaTunnelServer server = Mockito.mock(SeaTunnelServer.class);
+        CoordinatorService coordinatorService = 
newMockCoordinatorService(server);
+        CountDownLatch firstScheduleStarted = new CountDownLatch(1);
+        CountDownLatch allowFirstScheduleToFinish = new CountDownLatch(1);
+        try {
+            JobMaster blockedJobMaster =
+                    enqueueMockPendingJob(coordinatorService, 80001L, new 
CountDownLatch(1));
+            Mockito.when(blockedJobMaster.preApplyResources())
+                    .thenAnswer(
+                            invocation -> {
+                                firstScheduleStarted.countDown();
+                                while (!allowFirstScheduleToFinish.await(
+                                        100, TimeUnit.MILLISECONDS)) {
+                                    // Simulate a resource request that does 
not react immediately
+                                    // to master demotion or scheduler 
replacement.
+                                }
+                                return true;
+                            });
+
+            ReflectionUtils.setField(coordinatorService, "isActive", true);
+            invokePendingJobScheduler(coordinatorService);
+            Assertions.assertTrue(firstScheduleStarted.await(5, 
TimeUnit.SECONDS));
+
+            JobMaster secondJobMaster =
+                    enqueueMockPendingJob(coordinatorService, 80002L, new 
CountDownLatch(1));
+            invokePendingJobScheduler(coordinatorService);
+
+            await().atMost(5, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () -> {
+                                Mockito.verify(secondJobMaster, 
Mockito.times(1))
+                                        .preApplyResources();
+                                Assertions.assertFalse(
+                                        
coordinatorService.getPendingJobQueue().contains(80002L));
+                            });
+            Mockito.verify(blockedJobMaster, 
Mockito.times(1)).preApplyResources();
+        } finally {
+            allowFirstScheduleToFinish.countDown();
+            shutdownCoordinatorIfRunning(coordinatorService);
+        }
+    }
+
+    @Test
+    void testClearCoordinatorServiceDropsPendingJobsUnderRejectStrategy() 
throws Exception {
+        // Regression test for the duplicate-dispatch gap on the default 
REJECT schedule strategy:
+        // when the coordinator is cleared while a scheduler is blocked inside 
preApplyResources,
+        // the interrupted PendingJobInfo must NOT survive in pendingJobQueue, 
otherwise a
+        // restored scheduler (post same-node master flap-back) could 
re-dispatch the poisoned
+        // JobMaster after a fresh PendingJobInfo under the same job id has 
been enqueued.
+        SeaTunnelServer server = Mockito.mock(SeaTunnelServer.class);
+        EngineConfig engineConfig = new EngineConfig();
+        engineConfig.setScheduleStrategy(ScheduleStrategy.REJECT);
+        CoordinatorService coordinatorService = 
newMockCoordinatorService(server, engineConfig);
+        CountDownLatch allowFirstScheduleToFinish = new CountDownLatch(1);
+        try {
+            JobMaster blockedJobMaster =
+                    enqueueMockPendingJob(coordinatorService, 90001L, new 
CountDownLatch(1));
+            Mockito.when(blockedJobMaster.preApplyResources())
+                    .thenAnswer(
+                            invocation -> {
+                                allowFirstScheduleToFinish.await();
+                                return true;
+                            });
+
+            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());
+
+            // Simulate a master step-down. The blocked JobMaster is 
interrupted; the
+            // PendingJobInfo must be dropped from the queue so a later 
restore cannot
+            // re-dispatch the poisoned instance.
+            invokeClearCoordinatorService(coordinatorService);
+
+            await().atMost(5, TimeUnit.SECONDS)
+                    .untilAsserted(
+                            () -> {
+                                Assertions.assertFalse(
+                                        
coordinatorService.getPendingJobQueue().contains(90001L));
+                                Mockito.verify(blockedJobMaster, 
Mockito.atLeastOnce()).interrupt();
+                            });
+        } finally {
+            allowFirstScheduleToFinish.countDown();
+            shutdownCoordinatorIfRunning(coordinatorService);
+        }
+    }
+
     @Test
     void testPendingJobWithInsufficientResourceRespectsWaitStrategy() throws 
Exception {
         AtomicBoolean masterFlag = new AtomicBoolean(true);
@@ -548,6 +766,12 @@ public class CoordinatorServiceTest {
                 .ifPresent(ScheduledExecutorService::shutdownNow);
     }
 
+    private void invokePendingJobScheduler(CoordinatorService 
coordinatorService) throws Exception {
+        Method method = 
CoordinatorService.class.getDeclaredMethod("startPendingJobScheduleThread");
+        method.setAccessible(true);
+        method.invoke(coordinatorService);
+    }
+
     private JobMaster enqueueMockPendingJob(
             CoordinatorService coordinatorService, long jobId, CountDownLatch 
runLatch) {
         return enqueueMockPendingJob(coordinatorService, jobId, runLatch, 
true);
@@ -594,6 +818,12 @@ public class CoordinatorServiceTest {
                                         "Failed to get coordinator 
executorService by reflection"));
     }
 
+    private void shutdownCoordinatorIfRunning(CoordinatorService 
coordinatorService) {
+        if (!getCoordinatorExecutor(coordinatorService).isShutdown()) {
+            coordinatorService.shutdown();
+        }
+    }
+
     @Test
     public void 
testSeaTunnelEngineRetryableExceptionOperationCanBeRetryByHazelcast() {
 
@@ -1227,6 +1457,13 @@ public class CoordinatorServiceTest {
         method.invoke(coordinatorService);
     }
 
+    private void invokeClearCoordinatorService(CoordinatorService 
coordinatorService)
+            throws Exception {
+        Method method = 
CoordinatorService.class.getDeclaredMethod("clearCoordinatorService");
+        method.setAccessible(true);
+        method.invoke(coordinatorService);
+    }
+
     @Test
     public void testClearCoordinatorService() {
         JobInformation jobInformation =

Reply via email to