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

   ### Purpose of this pull request
   
   Adds one new Zeta E2E regression test covering the fix in #10448 for #10442
   ("seatunnel will never perform a checkpoint again once a previous checkpoint 
fails"), as part of
   the ongoing extreme-case E2E coverage initiative for the Zeta engine.
   
   **Issue #10442**: `CheckpointCoordinator#startTriggerPendingCheckpoint` 
increments `pendingCounter`
   before dispatching a checkpoint barrier. If that dispatch throws (the 
reporter's example:
   `CheckpointManager#sendOperationToMemberNode` -> 
`JobMaster#queryTaskGroupAddress` timing out under
   IMap contention), the old code just logged the exception and returned, 
leaving `pendingCounter`
   stuck above zero forever. Every later scheduled trigger attempt then saw 
`pendingCounter > 0` and
   just rescheduled itself, so the job kept reporting `RUNNING` with no error 
and no further
   checkpoints, silently, forever.
   
   **Fix #10448** ("[Fix][Zeta] make the job failed when triggering checkpoint 
fails (apache#10442)")
   replaces both catch blocks in that method with a call to 
`handleCoordinatorError(...,
   CheckpointCloseReason.CHECKPOINT_INSIDE_ERROR)`.
   
   ### Current `dev` behavior this test asserts (verified against the current 
HEAD, not assumed)
   
   I traced the full current call chain in `CheckpointCoordinator.java`,
   `CheckpointManager.java`, `JobMaster.java`, `SubPlan.java`, and 
`PhysicalPlan.java` before writing
   this test:
   
   - `startTriggerPendingCheckpoint`'s `catch (Exception e)` (around
     `CheckpointCoordinator.java:965-970`) calls
     `handleCoordinatorError("triggering checkpoint barrier failed", e, 
CheckpointCloseReason.CHECKPOINT_INSIDE_ERROR)`.
   - `handleCoordinatorError` (`CheckpointCoordinator.java:341-357`) marks the 
coordinator `FAILED`,
     calls `checkpointManager.handleCheckpointError(pipelineId, false)`, and 
resets `pendingCounter`
     to 0 via `cleanPendingCheckpoint`.
   - `CheckpointManager#handleCheckpointError` -> 
`JobMaster#handleCheckpointError`
     (`JobMaster.java:767-779`) -> `SubPlan#handleCheckpointError()` 
(`SubPlan.java:639-646`), which
     cancels the pipeline (`PipelineStatus.CANCELING`) if it is not already in 
an end state.
   - Once all tasks finish canceling, `SubPlan#getPipelineEndState()` 
(`SubPlan.java:243-290`) sees
     `canceledTaskNum > 0` and calls `cancelCheckpoint()`; because the 
coordinator's own
     `checkpointCoordinatorFuture` was already completed `FAILED` by 
`handleCoordinatorError`, that
     call returns the same `FAILED` status, which upgrades the pipeline's end 
state from `CANCELED`
     to `FAILED`.
   - With `job.retry.times = 0` (this test's conf), 
`SubPlan#canRestorePipeline()` is `false`
     (`getPipelineRestoreNum() < pipelineMaxRestoreNum` is `0 < 0`), so 
`stateProcess()`'s
     `FAILED`/`CANCELED` case skips restore and completes the pipeline future 
as `FAILED`.
   - `PhysicalPlan#addPipelineEndCallback` (`PhysicalPlan.java:148-196`) sees 
the (single) pipeline
     end as `FAILED`, so `failedPipelineNum > 0` and the whole job transitions 
to
     `JobStatus.FAILED` via `updateJobState`.
   
   So the current, documented behavior is: **the job reaches a terminal 
`FAILED` state**, not the
   old silent-forever-`RUNNING` hang. This is exactly what the merged fix's own 
unit test
   (`CheckpointBarrierTriggerErrorTest`) already asserts with mocks; this PR 
adds equivalent coverage
   at the E2E level with a real, unmocked, unmodified `queryTaskGroupAddress`.
   
   ### Trigger mechanism, and why it is reliable without having run it locally
   
   Per this project's Apache SeaTunnel local-verification policy, I have not 
run this test locally;
   GitHub CI on this PR's head is the first execution. I was correspondingly 
rigorous about tracing
   every step of the mechanism against the real, current source before writing 
the assertions:
   
   `CheckpointCoordinator#triggerCheckpoint` 
(`CheckpointCoordinator.java:1120-1132`) is the only code
   that can make `startTriggerPendingCheckpoint`'s 
`CompletableFuture.allOf(completableFutureArray)
   .get()` throw *synchronously* (as opposed to a per-task RPC merely failing 
later, which that
   `allOf` call never even notices, since it only waits for 
`triggerCheckpoint()` to return -- a
   separate, already-covered dead-letter scenario, see #12111). It maps every 
starting subtask through
   `checkpointManager::sendOperationToMemberNode` 
(`CheckpointManager.java:386-400`), which calls
   `jobMaster.queryTaskGroupAddress(...)` (`JobMaster.java:977-994`) *before* 
issuing the RPC. That
   method does exactly one thing that can throw: 
`ownedSlotProfilesIMap.get(pipelineLocation)`
   returning `null`, which throws `IllegalArgumentException("can't find task 
group address from
   taskGroupLocation: ...")`.
   
   I confirmed (via a repo-wide search) that `ownedSlotProfilesIMap`'s only 
entry-removal call site
   (`JobMaster#releasePipelineResource`, `JobMaster.java:923-950`) runs only 
after a pipeline has
   already left `RUNNING`, by which point `cleanPendingCheckpoint` has already 
cancelled that
   coordinator's own scheduler (`scheduler.shutdownNow()`, 
`CheckpointCoordinator.java:1203`) -- so
   nothing in the running system naturally races this lookup against a live, 
scheduled trigger.
   Killing/isolating a worker (this test class's usual technique) does not 
reach this code path
   either: a graceful leave fails the *task* directly via
   `CoordinatorService#failedTaskOnMemberRemoved` without ever touching this 
map, while an ungraceful
   one leaves a *stale but present* entry (the RPC fails later, asynchronously 
-- the separate
   dead-letter case covered by #12111, not this one).
   
   So this test uses a different, still entirely real, lever instead of cluster 
membership:
   `ownedSlotProfilesIMap` is a plain, named Hazelcast `IMap` 
(`Constant#IMAP_OWNED_SLOT_PROFILES`),
   obtained the exact same way this test class's own `getReadyToCloseCount` 
helper already reads
   `Constant#IMAP_RUNNING_JOB_STATE` directly, and the same way the 
engine-server module's own
   `EngineStateStoreMetricExportsTest` pokes this exact map in its unit tests. 
The test removes this
   job's entry from that live, shared map -- this is not a mock and not a 
reflected exception injected
   into production code: it is the same real, unmodified, running 
`queryTaskGroupAddress` that throws
   its own real `IllegalArgumentException` the next time it executes, exactly 
as it would if this
   bookkeeping ever went missing for any other reason. I checked every other 
reader of this map
   (metrics export, pipeline cleanup, 
`PhysicalVertex#checkTaskGroupIsExecuting` -- itself only
   reachable via master-failover restore, never steady-state `RUNNING`) and 
confirmed they all
   null-check and skip gracefully, so the removal cannot trip any other code 
path first.
   
   This is deterministic, not a narrow-window race like a worker kill: the 
entry is left removed
   permanently (the pipeline is about to fail anyway), so the very next 
scheduled trigger attempt
   that has not already started is guaranteed to observe the missing entry once 
the removal
   completes -- no timing window to miss.
   
   ### Test outline
   
   
`CheckpointCoordinatorFailoverIT#testStreamJobFailsAfterCheckpointTriggerDispatchFailure`:
   
   1. Starts a single embedded node running a `parallelism = 1` streaming 
FakeSource -> LocalFile job
      (`checkpoint.interval = 2000`, `job.retry.times = 0`).
   2. Waits for the job to be `RUNNING` and producing rows.
   3. Waits for the checkpoint-id counter to reach `2`, which (since 
`tryTriggerPendingCheckpoint`
      never allocates a new id while `pendingCounter > 0`) can only happen once 
checkpoint id 1 has
      fully completed -- proving checkpointing was healthy before the fault is 
injected.
   4. Removes the job's entry from `ownedSlotProfilesIMap` (the real-fault 
injection described above).
   5. Asserts the job reaches terminal `JobStatus.FAILED` within a bounded 
window, and that
      `JobResult#getError()` contains 
`CheckpointCloseReason.CHECKPOINT_INSIDE_ERROR`'s message text
      (traced end-to-end through `SubPlan`/`PhysicalPlan` back to the 
coordinator's own error).
   
   The bounded `Awaitility` wait for `FAILED` also implicitly demonstrates the 
pre-fix bug no longer
   occurs: against the old code, the job would stay `RUNNING` forever with no 
further checkpoints, so
   this wait would time out and fail the test.
   
   ### Does this PR introduce any user-facing change?
   
   No. Test-only; no `src/main` changes.
   
   ### How was this patch tested?
   
   - `./mvnw spotless:apply -pl 
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -nsu 
-Dmaven.gitcommitid.skip=true` -- BUILD SUCCESS, no further formatting changes 
needed.
   - Per this repo's Apache SeaTunnel local-verification policy, no local 
compile/test/E2E execution
     was performed; GitHub CI on this PR's head is the authoritative 
verification.
   
   ### Check list
   
   - [x] Test-only change, no `src/main` modifications.
   - [x] No new dependencies.
   - [x] No documentation changes needed (test-only).
   
   🤖 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