goutamadwant commented on code in PR #12314:
URL: https://github.com/apache/seatunnel/pull/12314#discussion_r4010811242


##########
seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointCoordinatorTest.java:
##########
@@ -1183,6 +1183,229 @@ void 
testRestoreCoordinatorShouldNotTriggerCheckpointWhenNotifyCompletedFails()
         }
     }
 
+    /**
+     * Builds an already-completeExceptionally'd {@code InvocationFuture} mock 
that behaves like a
+     * genuinely failed barrier-dispatch RPC for {@code 
CompletableFuture.allOf(...)}.
+     *
+     * <p>{@code InvocationFuture} is a final Hazelcast class with a 
package-private constructor, so
+     * it cannot be subclassed or instantiated directly from test code; it is 
mocked here using
+     * Mockito's inline mock maker (enabled for this module via {@code
+     * src/test/resources/mockito-extensions/org.mockito.plugins.MockMaker}), 
which mocks the class
+     * in place rather than generating a subclass. {@code 
completeExceptionally} is stubbed with
+     * {@code doCallRealMethod()} so the call genuinely runs the inherited 
{@code
+     * java.util.concurrent.CompletableFuture#completeExceptionally} 
implementation on this exact
+     * instance, setting its real internal completion state. That state is 
read directly (not
+     * through any overridable method) by the JDK's own {@code 
CompletableFuture.allOf(...)}
+     * internals, so the resulting future correctly observes this failure 
exactly as it would for a
+     * real, permanently-failed barrier RPC.
+     *
+     * @param cause the exception the simulated barrier dispatch RPC fails with
+     * @return a mock {@code InvocationFuture} already completed exceptionally 
with {@code cause}
+     */
+    private InvocationFuture<?> 
completedExceptionallyInvocationFuture(Throwable cause) {
+        InvocationFuture<?> future = Mockito.mock(InvocationFuture.class);
+        
Mockito.doCallRealMethod().when(future).completeExceptionally(Mockito.any(Throwable.class));

Review Comment:
   This helper does not create an exceptionally completed `InvocationFuture`. 
Mockito bypasses the Hazelcast constructor, leaving the future's internal state 
uninitialized, so calling the real `completeExceptionally` method here never 
completes the JDK stage. On Java 11, I reproduced the positive test timing out; 
the stale-callback test can pass without its callback running.
   
   Please use a real or faithfully initialized future and assert that both 
callback paths actually execute.



##########
seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointCoordinator.java:
##########
@@ -1016,6 +1031,48 @@ private void startTriggerPendingCheckpoint(
                                 CheckpointCloseReason.CHECKPOINT_INSIDE_ERROR);
                         return;
                     }
+                    // Observe the real per-task barrier dispatch futures 
asynchronously so a
+                    // dispatch failure is reported immediately with its real 
cause, instead of
+                    // being silently dropped until the checkpoint.timeout 
backstop below fires.
+                    // This must stay non-blocking: the checkpoint.timeout 
scheduled task is the
+                    // only backstop for a barrier RPC that never completes at 
all (member hang, or
+                    // a network partition before Hazelcast's own retries give 
up), so this method
+                    // must still return promptly and let that timeout get 
scheduled below.
+                    // Blocking here would also risk executor starvation: 
executorService is
+                    // ultimately the SynchronousQueue-backed pool created by
+                    // CoordinatorService#createCoordinatorExecutor, which 
rejects new work once
+                    // its threads are busy, so a blocking join on this thread 
could starve other
+                    // pipelines' checkpoint coordination sharing the same 
pool.
+                    CompletableFuture.allOf(barrierFutures)
+                            .whenCompleteAsync(
+                                    (ignored, throwable) -> {
+                                        if (throwable == null) {
+                                            return;
+                                        }
+                                        // Guard against a stale callback: by 
the time a slow RPC
+                                        // finally fails, this checkpoint may 
already have
+                                        // completed and been removed from 
pendingCheckpoints (see
+                                        // completePendingCheckpoint), or the 
whole coordinator may
+                                        // already be closed. 
handleCoordinatorError itself is a
+                                        // no-op once 
checkpointCoordinatorFuture is done, but
+                                        // pendingCheckpoints is keyed per 
checkpoint id, so check
+                                        // it explicitly to avoid failing the 
coordinator because
+                                        // of a late failure that belongs to 
an already completed
+                                        // checkpoint.
+                                        if (!pendingCheckpoints.containsKey(

Review Comment:
   Checking only map membership leaves a race after the last task 
acknowledgment. `PendingCheckpoint` becomes fully acknowledged and completes 
its future before `completePendingCheckpoint` persists the state, sends 
notifications, and removes this entry. If the barrier RPC was processed but its 
response later fails during that window, this callback fails the coordinator 
even though the checkpoint already succeeded.
   
   Please also require `!pendingCheckpoint.isFullyAcknowledged()`, as the 
timeout guard below already does, or make the state transition atomic and cover 
this race in a test.



##########
seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/checkpoint/CheckpointCoordinator.java:
##########
@@ -1016,6 +1031,48 @@ private void startTriggerPendingCheckpoint(
                                 CheckpointCloseReason.CHECKPOINT_INSIDE_ERROR);
                         return;
                     }
+                    // Observe the real per-task barrier dispatch futures 
asynchronously so a
+                    // dispatch failure is reported immediately with its real 
cause, instead of
+                    // being silently dropped until the checkpoint.timeout 
backstop below fires.
+                    // This must stay non-blocking: the checkpoint.timeout 
scheduled task is the
+                    // only backstop for a barrier RPC that never completes at 
all (member hang, or
+                    // a network partition before Hazelcast's own retries give 
up), so this method
+                    // must still return promptly and let that timeout get 
scheduled below.
+                    // Blocking here would also risk executor starvation: 
executorService is
+                    // ultimately the SynchronousQueue-backed pool created by
+                    // CoordinatorService#createCoordinatorExecutor, which 
rejects new work once
+                    // its threads are busy, so a blocking join on this thread 
could starve other
+                    // pipelines' checkpoint coordination sharing the same 
pool.
+                    CompletableFuture.allOf(barrierFutures)
+                            .whenCompleteAsync(

Review Comment:
   This callback can be rejected by the same coordinator executor. That pool 
uses a `SynchronousQueue` with `AbortPolicy`; when all coordinator threads are 
busy, `whenCompleteAsync` completes only its returned stage with 
`RejectedExecutionException` and this handler never runs. Because the returned 
stage is ignored, the dispatch failure is again delayed until the generic 
checkpoint timeout.
   
   Please use a callback path that cannot be dropped, or add an explicit 
rejection fallback and a saturated-pool regression test.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to