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]