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-12449-1882a579236c451df59ae8782a53c4460fb926f1
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit 146a1b5c510e1c33010a92981ae1afc8139c624b
Author: lw <[email protected]>
AuthorDate: Sun Sep 27 14:53:43 2026 +0000

    [Fix][E2E] Stabilize checkpoint trigger dispatch failure test (#12449)
    
    Co-authored-by: 201811510411lw <[email protected]>
---
 .../e2e/CheckpointCoordinatorFailoverIT.java       | 206 +++++++--------------
 1 file changed, 71 insertions(+), 135 deletions(-)

diff --git 
a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/CheckpointCoordinatorFailoverIT.java
 
b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/CheckpointCoordinatorFailoverIT.java
index 4d1aed50a6..d2256b3305 100644
--- 
a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/CheckpointCoordinatorFailoverIT.java
+++ 
b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/CheckpointCoordinatorFailoverIT.java
@@ -22,6 +22,7 @@ import 
org.apache.seatunnel.shade.org.apache.commons.lang3.tuple.ImmutablePair;
 import org.apache.seatunnel.common.config.Common;
 import org.apache.seatunnel.common.config.DeployMode;
 import org.apache.seatunnel.common.utils.FileUtils;
+import org.apache.seatunnel.common.utils.ReflectionUtils;
 import org.apache.seatunnel.engine.client.SeaTunnelClient;
 import org.apache.seatunnel.engine.client.job.ClientJobExecutionEnvironment;
 import org.apache.seatunnel.engine.client.job.ClientJobProxy;
@@ -617,108 +618,18 @@ public class CheckpointCoordinatorFailoverIT {
     }
 
     /**
-     * Regression test for the checkpoint-trigger-failure bug reported in <a
-     * href="https://github.com/apache/seatunnel/issues/10442";>#10442</a> and 
fixed by <a
-     * href="https://github.com/apache/seatunnel/pull/10448";>#10448</a> 
("[Fix][Zeta] make the job
-     * failed when triggering checkpoint fails (apache#10442)").
+     * Regression test for the checkpoint-trigger failure in <a
+     * href="https://github.com/apache/seatunnel/issues/10442";>#10442</a>, 
fixed by <a
+     * href="https://github.com/apache/seatunnel/pull/10448";>#10448</a>. A 
synchronous
+     * barrier-dispatch failure must fail the job through {@code 
CHECKPOINT_INSIDE_ERROR}, rather
+     * than leave it RUNNING forever with a pending checkpoint that prevents 
subsequent triggers.
      *
-     * <p>Before that fix, {@code 
CheckpointCoordinator#startTriggerPendingCheckpoint} (see {@code
-     * seatunnel-engine-server/.../checkpoint/CheckpointCoordinator.java} 
around lines 942-971)
-     * wrapped the checkpoint-barrier dispatch call like this:
-     *
-     * <pre>
-     * try {
-     *     CompletableFuture.allOf(completableFutureArray).get();
-     * } catch (InterruptedException e) {
-     *     throw new RuntimeException(e);
-     * } catch (Exception e) {
-     *     LOG.error(ExceptionUtils.getMessage(e));
-     *     return;
-     * }
-     * </pre>
-     *
-     * A {@code pendingCounter} field is incremented unconditionally right 
before this block ever
-     * runs (line ~1003, {@code pendingCounter.incrementAndGet();}) and is 
only ever decremented
-     * once a checkpoint fully completes (line ~1377). Before the fix, a 
dispatch failure here just
-     * logged and returned: {@code pendingCounter} stayed stuck above zero 
forever, and every later
-     * scheduled trigger attempt ({@code tryTriggerPendingCheckpoint}, line 
~800: {@code if
-     * (pendingCounter.get() > 0) { scheduleTriggerPendingCheckpoint(...); 
return; }}) would just
-     * reschedule itself and bail out without ever calling {@code 
createPendingCheckpoint} again.
-     * The job kept reporting {@code RUNNING} with no error and no further 
checkpoints, forever.
-     *
-     * <p>The fix (verified against the current {@code dev} HEAD before 
writing this test) replaces
-     * both catch blocks with a call to {@code handleCoordinatorError(..., 
CheckpointCloseReason
-     * .CHECKPOINT_INSIDE_ERROR)}, which marks the coordinator {@code FAILED}, 
calls {@code
-     * checkpointManager.handleCheckpointError(pipelineId, false)} (cancelling 
the pipeline via
-     * {@code SubPlan#handleCheckpointError()}), and resets {@code 
pendingCounter} to 0 as part of
-     * {@code cleanPendingCheckpoint}. Traced end to end for a single-pipeline 
job with restore
-     * disabled ({@code job.retry.times = 0}): {@code 
SubPlan#getPipelineEndState()} sees {@code
-     * canceledTaskNum > 0} and, because the checkpoint coordinator's own 
state is already {@code
-     * FAILED} by the time it calls {@code cancelCheckpoint()}, upgrades the 
pipeline's end state
-     * from {@code CANCELED} to {@code FAILED}; with restore disabled ({@code
-     * SubPlan#canRestorePipeline()} is false), {@code 
PhysicalPlan#addPipelineEndCallback} then
-     * fails the whole (single-pipeline) job. So the documented, current 
behavior this test asserts
-     * is: the job reaches a terminal {@code FAILED} state -- not 
silent-forever-{@code RUNNING}.
-     *
-     * <h2>Trigger mechanism</h2>
-     *
-     * <p>{@code CheckpointCoordinator#triggerCheckpoint} (line ~1120) is the 
only code that can
-     * make {@code startTriggerPendingCheckpoint}'s {@code 
CompletableFuture.allOf(...).get()} throw
-     * *synchronously*, as opposed to a per-task RPC merely failing later (a 
dead-letter scenario
-     * this same {@code allOf} bug never even notices, since it only waits for 
{@code
-     * triggerCheckpoint()} to return, not for the per-task futures inside its 
result to complete).
-     * {@code triggerCheckpoint} maps every starting subtask through {@code
-     * checkpointManager::sendOperationToMemberNode} 
(CheckpointManager.java:386-400), which calls
-     * {@code jobMaster.queryTaskGroupAddress(...)} (JobMaster.java:977-994) 
*before* issuing the
-     * RPC. That method does exactly one thing that can throw: {@code
-     * ownedSlotProfilesIMap.get(pipelineLocation)} returning {@code null}, 
which throws {@code
-     * IllegalArgumentException("can't find task group address from 
taskGroupLocation: ...")}.
-     *
-     * <p>A repo-wide search confirms {@code ownedSlotProfilesIMap}'s only 
entry-removal call site
-     * ({@code JobMaster#releasePipelineResource}, line ~949) runs only after 
a pipeline has
-     * *already* left {@code RUNNING}, by which point {@code 
cleanPendingCheckpoint} has already
-     * cancelled this coordinator's own scheduler (line ~1203, {@code 
scheduler.shutdownNow()}), so
-     * nothing in the running system naturally races this lookup against a 
live, scheduled trigger.
-     * Killing or isolating a worker -- this class's usual technique elsewhere 
-- does not help
-     * either: a graceful leave fails the *task* directly via {@code
-     * CoordinatorService#failedTaskOnMemberRemoved} without ever touching 
this map, while an
-     * ungraceful one leaves a *stale but present* entry (the RPC itself fails 
later, asynchronously
-     * -- exactly the dead-letter case {@code allOf} does not notice, and a 
different bug/test than
-     * this one).
-     *
-     * <p>So this test reaches for a different, still entirely real, lever 
instead of cluster
-     * membership: {@code ownedSlotProfilesIMap} is a plain, named Hazelcast 
{@code IMap} ({@code
-     * Constant#IMAP_OWNED_SLOT_PROFILES}), obtained the exact same way this 
class's own {@link
-     * #getReadyToCloseCount} already reads {@code 
Constant#IMAP_RUNNING_JOB_STATE} directly, and
-     * the same way the engine-server module's own {@code 
EngineStateStoreMetricExportsTest} pokes
-     * this exact map in its unit tests. Removing this job's entry from that 
live, shared map is not
-     * a mock and not a reflected exception injected into production code: it 
is the same real,
-     * unmodified, running {@code JobMaster#queryTaskGroupAddress} that throws 
its own real {@code
-     * IllegalArgumentException} the moment it next executes, exactly as it 
would if this
-     * bookkeeping ever went missing for any other reason. A check of every 
other reader of this map
-     * (metrics export, pipeline cleanup, {@code 
PhysicalVertex#checkTaskGroupIsExecuting} -- itself
-     * only reachable via master-failover restore, never during steady-state 
RUNNING) confirms all
-     * of them null-check and skip gracefully, so this removal cannot trip any 
other code path
-     * first.
-     *
-     * <p>This is deterministic, not a narrow-window race like a worker kill: 
the entry is left
-     * removed permanently (this pipeline is about to fail anyway), so unlike 
catching a kill at the
-     * exact moment a barrier is dispatched, the very next scheduled trigger 
attempt that has not
-     * already started -- or the one after that -- is guaranteed to observe 
the missing entry once
-     * the removal completes, with no timing window to miss. To also 
demonstrate the fault lands on
-     * a previously healthy coordinator (not one that was simply never able to 
checkpoint at all),
-     * the test first waits for the checkpoint-id counter to reach 2, which -- 
since {@code
-     * tryTriggerPendingCheckpoint} never allocates a new id while {@code 
pendingCounter > 0} (line
-     * ~800) -- can only happen after checkpoint id 1 has fully completed and 
been acknowledged.
-     *
-     * <p><b>What this test proves:</b> a real checkpoint-barrier dispatch 
failure, on a coordinator
-     * that was previously checkpointing successfully, fails the job (terminal 
{@code
-     * JobStatus.FAILED}, with an error message traceable to {@code 
CheckpointCloseReason
-     * #CHECKPOINT_INSIDE_ERROR}) within a bounded window. <b>What it 
implicitly also proves:</b>
-     * the pre-fix silent-forever-{@code RUNNING} behavior from #10442 no 
longer occurs -- had it,
-     * the bounded {@code Awaitility} wait below for {@code JobStatus.FAILED} 
would time out and
-     * fail this test, since the old code left the job {@code RUNNING} with no 
further checkpoints
-     * and no error, forever.
+     * <p>After a healthy checkpoint, remove this pipeline's real slot-profile 
bookkeeping so the
+     * next dispatch fails in {@code JobMaster#queryTaskGroupAddress}. 
Completion notifications use
+     * the same address lookup, so an allocated checkpoint ID alone is not a 
safe injection point.
+     * Wait for {@code pendingCounter == 0}, which is reached after completion 
notifications, and
+     * remove the entry while holding the coordinator's trigger lock. This 
excludes a new trigger
+     * between the idle check and removal without changing the production 
error-handling path.
      */
     @Test
     public void testStreamJobFailsAfterCheckpointTriggerDispatchFailure() 
throws Exception {
@@ -767,50 +678,67 @@ public class CheckpointCoordinatorFailoverIT {
                                         "Waiting for the source to start 
producing rows");
                             });
 
-            // Prove checkpointing is healthy before injecting the fault: the 
id counter can only
-            // reach 2 once checkpoint id 1 has been fully acknowledged -- see 
the class javadoc
-            // above for why (tryTriggerPendingCheckpoint never allocates a 
new id while
-            // pendingCounter is still above zero).
+            CheckpointCoordinator coordinator =
+                    getJobMaster(node, jobId)
+                            .getCheckpointManager()
+                            .getCheckpointCoordinator(pipelineId);
+            // Reflectively access the private lock and pendingCounter to time 
fault injection.
+            // Keep these field lookups in sync when refactoring 
CheckpointCoordinator.
+            Object triggerLock =
+                    ReflectionUtils.getField(coordinator, "lock")
+                            .orElseThrow(
+                                    () -> new IllegalStateException("Missing 
checkpoint lock"));
+            AtomicInteger pendingCounter =
+                    (AtomicInteger)
+                            ReflectionUtils.getField(coordinator, 
"pendingCounter")
+                                    .orElseThrow(
+                                            () ->
+                                                    new IllegalStateException(
+                                                            "Missing pending 
checkpoint counter"));
             CounterStateStore<String> checkpointCounterStore = 
checkpointCounterStore(node);
             String checkpointIdKey =
                     
StateStoreCheckpointIDCounter.convertLongIntToBase64(jobId, pipelineId);
+            IMap<PipelineLocation, Map<TaskGroupLocation, SlotProfile>> 
ownedSlotProfilesIMap =
+                    node.getMap(Constant.IMAP_OWNED_SLOT_PROFILES);
+            PipelineLocation pipelineLocation = new PipelineLocation(jobId, 
pipelineId);
             Awaitility.await()
                     .atMost(30, TimeUnit.SECONDS)
                     .pollInterval(200, TimeUnit.MILLISECONDS)
                     .untilAsserted(
                             () -> {
-                                Long currentId = 
checkpointCounterStore.get(checkpointIdKey);
-                                Assertions.assertNotNull(
-                                        currentId,
-                                        "waiting for the first checkpoint id 
to be allocated");
-                                Assertions.assertTrue(
-                                        currentId >= 2,
-                                        "waiting for checkpoint id 1 to be 
fully acknowledged"
-                                                + " before injecting the 
fault");
+                                Assertions.assertEquals(
+                                        JobStatus.RUNNING, 
clientJobProxy.getJobStatus());
+                                // Trigger creation and pendingCounter 
increment hold this lock.
+                                // Completion only decrements pendingCounter 
after notify succeeds.
+                                synchronized (triggerLock) {
+                                    Long currentId = 
checkpointCounterStore.get(checkpointIdKey);
+                                    Assertions.assertNotNull(
+                                            currentId,
+                                            "waiting for the first checkpoint 
id to be allocated");
+                                    Assertions.assertTrue(
+                                            currentId >= 2,
+                                            "waiting for at least one 
checkpoint to be triggered");
+                                    Assertions.assertEquals(
+                                            0,
+                                            pendingCounter.get(),
+                                            "waiting for checkpoint completion 
notifications"
+                                                    + " before injecting the 
dispatch failure");
+                                    Map<TaskGroupLocation, SlotProfile> 
removedSlotProfiles =
+                                            
ownedSlotProfilesIMap.remove(pipelineLocation);
+                                    Assertions.assertNotNull(
+                                            removedSlotProfiles,
+                                            "the running task's slot-profile 
bookkeeping should"
+                                                    + " exist before 
injection");
+                                    log.info(
+                                            "Job {} has no pending checkpoint; 
removed pipeline {}'s"
+                                                    + " slot-profile 
bookkeeping ({} task group(s))"
+                                                    + " before the next 
checkpoint-barrier dispatch.",
+                                            jobId,
+                                            pipelineId,
+                                            removedSlotProfiles.size());
+                                }
                             });
 
-            // Real-fault injection: remove this pipeline's entry from the 
same live, shared,
-            // named Hazelcast IMap (engine_ownedSlotProfilesIMap) that
-            // JobMaster#queryTaskGroupAddress consults on every 
checkpoint-barrier dispatch. See
-            // the class javadoc above for why this is real (not 
mocked/reflected),
-            // deterministic, and cannot be short-circuited by any other code 
path.
-            IMap<PipelineLocation, Map<TaskGroupLocation, SlotProfile>> 
ownedSlotProfilesIMap =
-                    node.getMap(Constant.IMAP_OWNED_SLOT_PROFILES);
-            PipelineLocation pipelineLocation = new PipelineLocation(jobId, 
pipelineId);
-            Map<TaskGroupLocation, SlotProfile> removedSlotProfiles =
-                    ownedSlotProfilesIMap.remove(pipelineLocation);
-            Assertions.assertNotNull(
-                    removedSlotProfiles,
-                    "the running task's slot-profile bookkeeping should exist 
before injection");
-            log.info(
-                    "Job {} checkpoint id counter reached 2; removed pipeline 
{}'s slot-profile"
-                            + " bookkeeping ({} task group(s)) so the next 
checkpoint-barrier"
-                            + " dispatch hits CheckpointCoordinator's real, 
unmodified"
-                            + " queryTaskGroupAddress failure path.",
-                    jobId,
-                    pipelineId,
-                    removedSlotProfiles.size());
-
             Awaitility.await()
                     .atMost(60, TimeUnit.SECONDS)
                     .pollInterval(500, TimeUnit.MILLISECONDS)
@@ -832,6 +760,14 @@ public class CheckpointCoordinatorFailoverIT {
                                     + " coordinator's CHECKPOINT_INSIDE_ERROR 
path (see"
                                     + " 
CheckpointCoordinator#handleCoordinatorError), but got: "
                                     + jobResult.getError());
+            Assertions.assertTrue(
+                    jobResult.getError().contains("can't find task group 
address"),
+                    () ->
+                            "Expected the injected address lookup failure, but 
got: "
+                                    + jobResult.getError());
+            Assertions.assertTrue(
+                    
jobResult.getError().contains("CheckpointCoordinator.triggerCheckpoint"),
+                    () -> "Expected a barrier-dispatch failure, but got: " + 
jobResult.getError());
         } finally {
             if (engineClient != null) {
                 engineClient.close();

Reply via email to