davidzollo opened a new pull request, #12033:
URL: https://github.com/apache/seatunnel/pull/12033
## Purpose
Adds multi-node E2E regression coverage for the CANCELING-stuck-forever race
fixed in #10729 ("[Fix][Zeta] prevent cancel stuck and downgrade tmp cleanup
failure", commit `2c81f7f3c31`).
Before that fix, `PhysicalVertex#noticeTaskExecutionServiceCancel` sent a
`CancelTaskOperation` to the worker owning the task and retried while that
worker remained a cluster member, but once the worker actually left the
cluster before acking the cancel, the retry loop simply exited without ever
resolving the task's state. The vertex — and therefore its pipeline and job
— stayed `CANCELING` forever, because nothing else in that method drove it
to a terminal state. The fix added a fallback right after the retry loop: if
the loop exits with the ack still missing and the vertex is still
`CANCELING`, mark it `CANCELED` locally.
The fix shipped with no dedicated regression test reproducing the original
multi-node timing/race at the engine level.
## What was still missing
`CoordinatorService#failedTaskOnMemberRemoved` also matches `CANCELING`
tasks when a member is lost (this branch predates the fix — confirmed via
`git show` against the fix commit's parent), which raises the question of
whether that path already covered the bug. It does not, reliably: it has to
go through the same `synchronized updateTaskState` as
`PhysicalVertex#cancel`/`#updateTaskState`/`#stateProcess`, all of which are
synchronized on the vertex itself, and the whole call chain from `cancel()`
down into `noticeTaskExecutionServiceCancel` runs on that lock without
releasing it. A competing `failedTaskOnMemberRemoved` call can therefore
only run *after* the cancelling thread has already returned — by which
point the fixed code has already resolved the vertex to `CANCELED`, and the
end-state consistency check inside `updateTaskState` rejects any later
attempt to move a terminal-state task to `FAILED`. So the fixed method's own
fallback is the sole, deterministic mechanism that unsticks this specific
race, which is why the new test asserts `CANCELED` (not `FAILED`) as the
outcome.
I also checked the existing fault-tolerance suite
(`ClusterFaultToleranceIT`, `SplitClusterFaultToleranceIT`,
`ClusterFailureNoRestoreIT`) for anything that already combines a cancel
request with a worker crash. It does not: existing tests either cancel a
job on a fully healthy cluster, or kill a worker on a job that is simply
`RUNNING` with no cancel in flight and cancel only after restore completes.
None land a worker crash inside the narrow window between "cancel issued,
`CANCELING` set" and "cancel ack received".
## What this test does
`SplitClusterFaultToleranceIT#testStreamJobCancelResolvesWhenWorkerCrashesBeforeCancelAck`:
1. Starts a minimal split-role cluster: 1 dedicated master + 1 dedicated
worker.
2. Submits the same long-running streaming holder job used elsewhere in
this suite (`pending_jobs_streaming_lifecycle.conf`) and waits for the
top-level job status **and** every individual task vertex to report
`RUNNING`, so the cancel below isn't confounded by an in-flight
`DEPLOYING` transition.
3. Issues `cancelJob()` on its own thread. This is required because
`CoordinatorService#cancelJob` runs `JobMaster#cancelJob` synchronously
on its own executor, and the client's `cancelJob()` call blocks
(`PassiveCompletableFuture#join`) until that whole per-vertex cancel
resolves — so the test thread needs to stay free to react while the
cancel is in flight on the server.
4. Tight-polls (busy-spin, not Awaitility's ~100ms default interval) every
task vertex's `PhysicalVertex#getExecutionState()` via the same in-JVM
white-box access as
`SplitClusterPendingJobLifecycleFailoverIT#getJobMaster`,
and the instant any vertex is first observed `CANCELING`, immediately
shuts the worker down.
5. Asserts the job converges to `CANCELED` within a bounded time (using an
already-in-flight `waitForJobCompleteV2()` future so the assertion
itself cannot hang forever if the regression reappears), and that the
original `cancelJob()` call also returns without hanging.
## Why this trigger-construction is reliable
The race window this bug depends on — between a vertex's state flipping to
`CANCELING` and the `CancelTaskOperation` ack coming back from the worker —
can be sub-millisecond on a healthy local cluster, so blind timing (sleep
then kill) would be unreliable. Instead the test uses direct white-box
polling of internal engine state (same technique already proven in this
series, e.g. #12027), reacting the instant the state transition is
observed rather than guessing a delay. Because `SubPlan#cancelPipeline`
cancels vertices sequentially on a single thread and each vertex's cancel
call blocks synchronously on `noticeTaskExecutionServiceCancel`, the
`CANCELING` state is guaranteed to be visible for at least the duration of
a real cross-process Hazelcast operation invocation (these test nodes
communicate over real loopback sockets, not an in-memory bypass) — long
enough for a tight busy-poll loop to reliably observe it before the ack
round-trip completes.
## Test plan
- New test only; no production code changed.
- `./mvnw spotless:apply` run on the affected module — no additional
formatting changes beyond what was already applied to the new code.
- `./mvnw install -DskipTests` run on the affected module and its
dependency chain (this repo's shaded modules need `install`, not bare
`compile`/`test-compile`, to resolve relocated classes) — confirmed
`BUILD SUCCESS`, and confirmed the new test class actually compiled
(`target/test-classes` contains the compiled `.class` file for the new
method) rather than relying on a build that skips test compilation
entirely.
- Full test execution (real Hazelcast cluster startup, worker crash
injection) is left to CI per this repository's E2E conventions; not run
in this sandbox.
--
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]