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]

Reply via email to