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

   ### Purpose
   
   `CoordinatorExecutorMassFailoverStormIT` (added by #12104) submits 30
   concurrent streaming jobs, kills the active master mid-flight to force a
   mass-simultaneous restore, asserts the standby master's shared coordinator
   executor grows past its configured core size, and then tears down by
   cancelling every job and awaiting each job's own terminal `JobStatus`. The
   teardown asserted that status was `CANCELED` or `FINISHED`. That assumption
   is wrong: a job cancelled while this storm is in progress can legitimately
   end `FAILED`.
   
   ### The failure
   
   Seen on an unrelated PR's CI run (`zeta-bounded-stream-failover-completion`,
   run `34150315578`, job `engine-v2-it (8, ubuntu-latest)`, attempt 1, job id
   `101831579495`):
   
   ```
   org.opentest4j.AssertionFailedError: Job 1149424434122522626 should have
   reached a terminal status, was FAILED ==> expected: <true> but was: <false>
        at 
CoordinatorExecutorMassFailoverStormIT.cancelAllAndAwaitTerminal(CoordinatorExecutorMassFailoverStormIT.java:399)
   ```
   
   ### Log-evidenced mechanism
   
   Job `1149424434122522626` ("mass_retry_job_7") reached `FAILED` instead of
   `CANCELED` through the following chain, confirmed both from the raw CI log
   and by reading the corresponding engine source on `dev`:
   
   1. **Checkpoint expiry during the storm.** This module's `seatunnel.yaml`
      sets `checkpoint.timeout=100000` (100s). The job's first checkpoint went
      unacknowledged for that whole window while the storm was in progress. At
      `19:24:03,609` the checkpoint's own pre-scheduled expiry check (the
      `scheduler.schedule(...)` task started inside
      `CheckpointCoordinator#triggerCheckpoint`,
      `CheckpointCoordinator.java:986-996`) fired:
      ```
      ERROR CheckpointCoordinator - trigger checkpoint failed
      CheckpointException: Checkpoint expired before completing. Please increase
      checkpoint timeout in the seatunnel.yaml or jobConfig env.
      ```
      This aborts the checkpoint and marks the `CheckpointCoordinator` `FAILED`,
      which in turn tells every task in the pipeline to cancel
      (`SubPlan` logs `Pipeline: [(1/1)] turned from state RUNNING to CANCELING`
      at `19:24:03,526`).
   
   2. **One task's cleanup RPC does not fail fast.** Three of the pipeline's
      four tasks reach `CANCELED` within about a second. The fourth
      (`SourceTask (1/2)`, `taskGroupId=3`) was, at that same instant, mid-RPC
      (`ReportMetricsOperation`) to the master this test independently kills.
      Its connection closes immediately at `19:24:03,752`
      (`TcpServerConnection ... closed. Reason: Exception in Connection`), but
      the pending Hazelcast invocation does not fail fast: it only surfaces once
      the `operation-heartbeat-timeout` gives up on it, `120298`ms later, at
      `19:26:04,011`:
      ```
      ERROR TaskExecutionService - failed to update metrics, 
payloadTaskCount=100,
      invocationLatencyMs=120298.
      OperationTimeoutException: ReportMetricsOperation invocation failed to
      complete due to operation-heartbeat-timeout. ... target=[localhost]:5802
      ```
      (`TaskExecutionService#updateMetricsContextInImap`,
      `TaskExecutionService.java:880` and `893-901`.) The task still completes
      `CANCELED` immediately afterward, since it was already flagged for
      cancellation — just ~120s after its three siblings.
   
   3. **The pipeline's final state re-checks the checkpoint coordinator.** Once
      all four tasks have completed, `SubPlan#getPipelineEndState()`
      (`SubPlan.java:250-259`) makes the call: it sees every task ended
      `CANCELED` (so its default verdict would be `CANCELED`), but this method
      also re-checks the `CheckpointCoordinator`'s own status. Since that was
      marked `FAILED` in step 1, it overrides the pipeline's terminal state to
      `FAILED`:
      ```
      INFO SubPlan - Job mass_retry_job_7 (1149424434122522626), Pipeline: 
[(1/1)]
      will end with state FAILED
      ```
      The job then follows `CANCELING -> FAILING -> FAILED`, and the test's
      teardown, which only accepted `CANCELED`/`FINISHED`, fails.
   
   ### Why the old predicate was wrong and the new one is right
   
   The teardown's actual job is only to prove that cancelling all 30 jobs does
   not leave any of them stuck in a non-terminal state after the mass-failover
   storm — it is cleanup, not the property under test. The property under test
   is the coordinator executor's growth, asserted earlier in the same method
   and **left completely untouched by this PR**. `CANCELED`/`FINISHED` was an
   incomplete enumeration of "terminal"; `JobStatus#isEndState()` is the
   engine's own definition (it also covers `SAVEPOINT_DONE` and `UNKNOWABLE`)
   and is the correct predicate for "did not leak past the test", regardless of
   which specific terminal state a job under storm conditions happens to land
   in. The per-job assertion message is unchanged, and the per-status outcome
   counts are now also logged at INFO so CI output keeps visibility into how
   many jobs ended `CANCELED` vs `FAILED` on any given run.
   
   ### Engine observation (not asserted here)
   
   Two things surfaced during this investigation that are not fixed or
   asserted on by this PR, but may be worth a maintainer's look:
   
   - A task's own status/metrics RPC to a master whose connection has already
     closed does not fail fast; it blocks for the full ~120s Hazelcast
     `operation-heartbeat-timeout` before the invocation surfaces as failed.
     That is up to ~120s of avoidable latency on every affected task's own
     terminal transition during a master failover.
   - Separately, `SubPlan#getPipelineEndState()` decides a pipeline whose tasks
     all individually ended `CANCELED` should nonetheless be reported `FAILED`
     whenever the checkpoint coordinator's own last status was `FAILED` — even
     though the checkpoint failure is what triggered the cancellation in the
     first place. That may be intentional (a checkpoint-driven cancel is not a
     "clean" cancel), but it means a client-observed `CANCELED` cancel request
     can surface to the user as `FAILED` purely because of checkpoint timing
     during a master switch. Worth a maintainer's look as a possible
     fast-fail / terminal-state-mapping improvement.
   
   ### Scope
   
   Test-only change to
   `CoordinatorExecutorMassFailoverStormIT#cancelAllAndAwaitTerminal`. No
   production code, no `pom.xml`, and none of the executor-growth assertions
   (the actual property under test) are touched. `git diff --stat` for this PR
   shows exactly one file changed.
   
   ### Validation
   
   - `./mvnw spotless:apply -pl 
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -nsu 
-Dmaven.gitcommitid.skip=true` — no additional formatting changes.
   - No local build/test was run for this change, per current SeaTunnel
     contribution guidance for this repository; correctness is verified through
     GitHub CI on this PR's head commit.
   
   🤖 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