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

   ## What this PR verifies and covers
   
   `CheckpointCoordinator#startTriggerPendingCheckpoint` triggers a checkpoint 
barrier like this
   
(`seatunnel-engine/seatunnel-engine-server/.../checkpoint/CheckpointCoordinator.java:944-958`):
   
   ```java
   CompletableFuture<InvocationFuture<?>[]> completableFutureArray =
           CompletableFuture.supplyAsync(() -> new CheckpointBarrier(...), 
executorService)
                   .thenApplyAsync(this::triggerCheckpoint, executorService);
   CompletableFuture.allOf(completableFutureArray).get();
   ```
   
   `completableFutureArray` is a single `CompletableFuture` whose *value* is an
   `InvocationFuture<?>[]` -- not the array itself. 
`CompletableFuture.allOf(...)` only spreads a
   real array argument into its varargs; handed one future, it waits on exactly 
that future, i.e.
   until `triggerCheckpoint()` returns (the per-task 
`CheckpointBarrierTriggerOperation` RPCs are
   *fired*), not until any of them land or are acknowledged. The individual 
`InvocationFuture`s
   inside the resolved array are never consulted again by anything in this 
class, so a barrier-RPC
   failure to a since-unreachable worker is silently discarded -- a true dead 
letter, not merely
   delayed handling.
   
   Contrast the same class's own correct pattern for `notifyTaskStart()` / 
`notifyCompleted()`
   (lines ~451-452, ~474-481): both spread an already-resolved 
`InvocationFuture<?>[]` directly into
   `allOf`, which genuinely waits on every element.
   
   I independently re-verified this against the current `dev` HEAD before 
writing the test --
   confirmed the exact code, line numbers, and contrasting patterns cited above 
are still present
   and unchanged.
   
   The only backstop for a barrier-dispatch RPC that never lands is the 
scheduled per-pending-checkpoint
   timeout later in the same method (lines ~972-1001): once 
`checkpoint.timeout` elapses without the
   checkpoint becoming fully acknowledged, 
`CheckpointCloseReason#CHECKPOINT_EXPIRED` fires, cancelling
   and restarting the pipeline the same way a hard task failure would.
   
   ## Test added
   
   
`CheckpointCoordinatorFailoverIT#testStreamJobRecoversAfterWorkerUnreachableDuringCheckpointBarrierDispatch`
   
   - Starts a 1-master/2-worker cluster and runs a parallelism=1 streaming job 
(single pipeline,
     single task) so the target worker can be identified unambiguously.
   - Identifies which worker hosts the task via its live execution address, 
then tightly polls
     (10ms) the shared checkpoint-id counter state store used elsewhere in this 
file, terminating the
     target worker in the same loop iteration that first observes a new 
checkpoint id -- i.e. as soon
     as that checkpoint's barrier dispatch is imminent or already starting. 
Termination uses
     `HazelcastInstance.getLifecycleService().terminate()` (not `shutdown()`), 
which per Hazelcast's
     own `Node#shutdown(boolean terminate)` skips the graceful cluster-leave 
notice that would
     otherwise let `CoordinatorService#failedTaskOnMemberRemoved` fail the task 
almost immediately --
     that fast, generic path is exactly what the existing 
`ClusterFaultToleranceIT`-style worker-kill
     tests exercise, and it would mask this specific dispatch-window bug 
entirely.
   - The job's `checkpoint.timeout` (8s) is deliberately far below Hazelcast's 
own native membership
     failure-detection ceiling. That ceiling is *not* incidental here: this 
build's shaded Hazelcast
     jar compiles in a 60-second default for 
`hazelcast.max.no.heartbeat.seconds` (verified by
     disassembling `com.hazelcast.spi.properties.ClusterProperty`'s static 
initializer), and this
     test module's own `hazelcast.yaml` does not raise it -- only this repo's 
top-level,
     production-only `config/hazelcast.yaml` does, and that file is not on this 
module's test
     classpath. Rather than depend on whichever `hazelcast.yaml` happens to be 
on the classpath, the
     test explicitly overrides `hazelcast.max.no.heartbeat.seconds` (to 180s) 
on every node it starts,
     so Hazelcast's own native failure detection provably cannot fire inside 
the test's bounded
     recovery wait, regardless of ambient config now or in the future.
   - Asserts the job recovers: status back to RUNNING, the task redeployed onto 
the surviving worker
     (by execution address), and row output resuming -- all within a bounded 
window tied to
     `checkpoint.timeout`, not Hazelcast's much slower native detection.
   
   **What this test proves:** the job recovers within a bounded time tied to 
`checkpoint.timeout`,
   even though the coordinator cannot detect this specific barrier-dispatch RPC 
failure immediately.
   
   **What it does NOT prove:** that the barrier-dispatch RPC failure is caught 
immediately -- per the
   bug above, it is not, and this test does not assert instant detection. It 
also does not modify
   `CheckpointCoordinator` itself; this is test-only coverage of 
currently-fragile behavior.
   
   ## Test plan
   
   - `./mvnw spotless:apply -pl 
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` -- 
passed.
   - Compilation verified locally: `./mvnw install -pl 
'seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base' -nsu 
-Dmaven.gitcommitid.skip=true -DskipTests -Dspotless.check.skip=true` -- BUILD 
SUCCESS, with the compiled `CheckpointCoordinatorFailoverIT.class` (containing 
the new test method) confirmed present under `target/test-classes`. This 
narrower, non-`-am` invocation resolved all upstream dependencies from a local 
repository already freshly built from this same unmodified `dev` base (this 
machine was running many concurrent, unrelated builds during verification); a 
full `-am` invocation from a clean state is expected to succeed identically 
since no `src/main/**` code was touched.
   - Actual test execution: per this repo's Apache SeaTunnel local-verification 
policy, this PR relies on GitHub CI (the fork's Actions run behind the `Build` 
check) as the authoritative execution result; local test execution was not 
performed.
   
   🤖 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