davidzollo opened a new pull request, #12314:
URL: https://github.com/apache/seatunnel/pull/12314

   ## Summary
   
   `CheckpointCoordinator#startTriggerPendingCheckpoint` fires every per-task 
checkpoint barrier RPC via `triggerCheckpoint()` (which returns 
`InvocationFuture<?>[]`, `CheckpointCoordinator.java:1167-1179`), but only ever 
waited for:
   
   ```java
   CompletableFuture.allOf(completableFutureArray).get();
   ```
   
   `completableFutureArray` is `CompletableFuture<InvocationFuture<?>[]>` - the 
future that resolves once `triggerCheckpoint()` has *produced* the array, i.e. 
once every barrier RPC has been *sent* via `sendOperationToMemberNode`. Because 
`InvocationFuture` extends the JDK `java.util.concurrent.CompletableFuture` 
(`InvocationFuture -> AbstractInvocationFuture -> InternalCompletableFuture -> 
java.util.concurrent.CompletableFuture`), the type system silently accepted 
wrapping that single array-producing future as one varargs element to 
`allOf(...)`. The per-task RPCs themselves - the actual `InvocationFuture` 
elements inside the array - were never observed anywhere else in the method 
(the array is used only to build this one `allOf(...)` call).
   
   Consequence: a barrier RPC that fails permanently (target member left, or a 
non-retryable Hazelcast operation failure) is a dead letter. The only backstops 
are the `checkpoint.timeout` scheduled task 
(`CheckpointCoordinator.java:1019-1047` after this change, previously 
~`:972-1001`) which reports a generic `CHECKPOINT_EXPIRED` only after the full 
timeout elapses, and member-removal handling. This is bounded (not P0) but 
delays detection and hides the real failure cause behind a generic timeout.
   
   This was raised during review of #12111 (E2E coverage for this exact 
dead-letter scenario); the reviewer requested a production fix rather than 
test-only coverage.
   
   ## Fix
   
   In `startTriggerPendingCheckpoint` (the only method touched, 
`CheckpointCoordinator.java`):
   
   1. Keep the existing `.get()` on `completableFutureArray` unchanged in 
behavior (same catch blocks, same `CHECKPOINT_INSIDE_ERROR` handling for the 
synchronous-throw path used by #10448) - now capturing the returned array into 
`InvocationFuture<?>[] barrierFutures`.
   2. Observe the real per-task futures **asynchronously and non-blocking**: 
`CompletableFuture.allOf(barrierFutures).whenCompleteAsync((ignored, throwable) 
-> {...}, executorService)`.
      - Non-blocking is essential: the `checkpoint.timeout` scheduled task 
right after this block is the only backstop for an RPC that never completes at 
all (member hang, network partition before Hazelcast's own retries give up); a 
blocking join here would prevent that timeout from ever getting scheduled.
      - `executorService` is ultimately the `SynchronousQueue`-backed pool 
created by `CoordinatorService#createCoordinatorExecutor` 
(`CoordinatorService.java:287-298`), which rejects new work once its threads 
are busy - a blocking wait on this thread could starve other pipelines' 
checkpoint coordination sharing the same pool.
   3. On a real failure, guard against a stale callback: by the time a slow RPC 
finally fails, the checkpoint may already have completed and been removed from 
`pendingCheckpoints` (`completePendingCheckpoint`, 
`CheckpointCoordinator.java:1420`), or the whole coordinator may already be 
closed. The callback checks `pendingCheckpoints.containsKey(checkpointId)` 
(mirroring the same idiom already used by the `checkpoint.timeout` task) before 
calling `handleCoordinatorError`, so a late failure belonging to an 
already-completed checkpoint is a no-op.
   4. Route a real failure to `handleCoordinatorError(..., 
CheckpointCloseReason.CHECKPOINT_INSIDE_ERROR)` - the same reason already used 
for the synchronous-throw path in this method (the #10448 path), so the 
coordinator now reports the actual dispatch failure immediately instead of a 
generic `CHECKPOINT_EXPIRED` after the full timeout.
   
   Blast radius: every checkpoint and savepoint trigger goes through this 
method, so this changes the failure-detection latency/cause for *any* barrier 
dispatch RPC failure across all pipelines - previously silent-until-timeout, 
now immediate with the real cause. No behavior changes when barrier RPCs 
succeed.
   
   ## Test-double fix (required by this change)
   
   `CheckpointCoordinatorTest`'s `TestCheckpointManager` and a couple of inline 
`Mockito.mock(CheckpointManager.class)` stubs (`buildMinimalCoordinator`) 
return `null` from `sendOperationToMemberNode(...)`. That `null` previously 
never reached `allOf(...)` since only the array-producing future was observed; 
with this fix, `CompletableFuture.allOf(barrierFutures)` would 
`NullPointerException` synchronously on a `null` array element, and since that 
call sits inside a `thenAccept(...)` continuation whose result is never read, 
the exception would be silently swallowed - hanging the checkpoint forever with 
no error surfaced.
   
   I audited every existing test that reaches `startTriggerPendingCheckpoint` 
for real (not via a `Mockito.doNothing()`-stubbed 
`tryTriggerPendingCheckpoint`):
   - `testFilteringClosedTasksAndActions` uses `TestCheckpointManager` (returns 
`null`), but calls `coordinator.triggerCheckpoint(barrier)` **directly** - it 
never goes through `startTriggerPendingCheckpoint`, so it's unaffected.
   - `testSchedulerThreadShouldNotBeInterruptedBeforeJobMasterCleaned`, 
`testCheckpointContinuesWorkAfterClockDrift`, `testCheckpointMinPause` all use 
a real (or real-delegating) `CheckpointManager` whose 
`sendOperationToMemberNode` dispatches a genuine (non-null) Hazelcast 
`InvocationFuture`.
   - No existing test's happy path actually feeds a `null` element into the new 
`allOf(barrierFutures)` call. No test-double changes were needed to keep 
existing tests passing.
   
   ## New tests (`CheckpointCoordinatorTest.java`)
   
   - **`testBarrierDispatchFailureIsRoutedToCoordinatorError`**: stubs 
`sendOperationToMemberNode` to return an already `completeExceptionally`'d 
`InvocationFuture`. Since `InvocationFuture` is `final` with a package-private 
constructor (cannot be subclassed or instantiated from test code), it is mocked 
via Mockito's inline mock maker (already enabled for this module via 
`src/test/resources/mockito-extensions/org.mockito.plugins.MockMaker` = 
`mock-maker-inline`, and already used elsewhere in this file for the same class 
and for `Mockito.mockStatic`). `completeExceptionally` is stubbed with 
`doCallRealMethod()` so it genuinely runs the inherited JDK 
`CompletableFuture#completeExceptionally` on that instance, setting its real 
internal completion state - which `CompletableFuture.allOf(...)`'s own 
internals read directly. The checkpoint timeout is configured to 10 minutes, 
far longer than the test's 10-second Awaitility bound, so a pass can only be 
explained by the new asynchronous path,
  not the timeout backstop. Asserts `handleCoordinatorError` is invoked with 
`CHECKPOINT_INSIDE_ERROR` and the real cause (allowing for the 
`CompletionException` wrapping that `CompletableFuture.allOf(...)` applies per 
its documented contract).
   - **`testStaleBarrierDispatchFailureForCompletedCheckpointIsIgnored`**: same 
failing RPC, but the checkpoint id is deliberately never registered in 
`pendingCheckpoints` (simulating a checkpoint that already completed and was 
removed). Asserts `handleCoordinatorError` is never called within a bounded 
wait - the stale-callback guard.
   
   ## Verification
   
   - `./mvnw spotless:apply -pl seatunnel-engine/seatunnel-engine-server -nsu 
-Dmaven.gitcommitid.skip=true` - clean, no additional formatting changes.
   - Per this task's constraints, no local compile, test, or other build/run 
verification was performed. Every symbol, signature, and generic type used in 
the production change and the new tests was manually verified against the 
pinned source at this PR's base commit (including a `javap` inspection of the 
shaded Hazelcast jar to confirm `InvocationFuture`'s exact class hierarchy and 
that it is `final`). GitHub CI on this PR's head is the verification of record 
for compilation and test execution.
   
   ## Test plan
   
   - [ ] GitHub CI (compile + `seatunnel-engine-server` unit tests) passes on 
this PR's head.
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)
   


-- 
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