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

   ### Purpose of this pull request
   
   Closes the remaining gap in #10834 ("Job stuck permanently after master 
failover, unable to complete") that #10836 did not actually cover.
   
   #10834 named three affected scenarios; #10836 fixed one root cause 
(`readyToCloseStartingTask` not surviving failover) and added 
`CheckpointCoordinatorFailoverIT` tests for it. But those tests only cover:
   - `testBatchJobCompletesAfterMasterFailover`: BATCH natural completion after 
failover.
   - `testStreamJobContinuesAfterMasterFailover`: streaming checkpoint-id 
continuity after failover (job is manually `cancelJob()`'d at the end, never 
actually completes).
   
   Neither exercises the third scenario #10834's own text calls out: *"a 
Streaming job triggers a cancel or savepoint and master failover happens during 
the shutdown phase."* Investigating this gap (as part of a broader initiative 
auditing Zeta failover edge cases this week) found that the **savepoint** 
variant is a real, currently-unfixed regression, confirmed both by reading the 
source and by reproducing it with a new E2E test before the fix was in place.
   
   ### Root cause
   
   - `JobMaster#savePoint()` sets `JobStatus#DOING_SAVEPOINT` and triggers a 
`SAVEPOINT_TYPE` checkpoint via 
`CheckpointManager#triggerSavePointsAndWaitComplete()`.
   - If the master dies before that checkpoint is fully acknowledged, 
`CoordinatorService#restoreJobFromMasterActiveSwitch` unconditionally calls 
`PhysicalPlan#updateJobState(JobStatus.PENDING)`. 
`PhysicalPlan#stateProcess()`'s `PENDING` case cascades this straight through 
to `RUNNING` — nothing anywhere in that path checks whether the job's status 
was `DOING_SAVEPOINT` before the failover, so that marker is silently lost.
   - `CheckpointCoordinator#restoreCoordinator` independently discards the 
in-flight `SAVEPOINT_TYPE` `PendingCheckpoint` via 
`cleanPendingCheckpoint(CHECKPOINT_COORDINATOR_RESET)`, and once the pipeline's 
tasks are confirmed running again, re-triggers only a regular `CHECKPOINT_TYPE` 
checkpoint. Nothing re-derives "a savepoint was requested" from anywhere, and 
`startSavepoint()` is never called again.
   
   Net effect: the client's savepoint request is silently and permanently 
dropped. The job keeps running and periodically checkpointing as if the 
savepoint had never been requested — exactly the "stuck permanently, unable to 
complete" signature #10834 describes, for the one scenario its own fix and 
tests never actually exercised.
   
   This is unlike the plain `cancelJob()` path (already covered by 
`testStreamJobContinuesAfterMasterFailover`): `SubPlan` has a real 
`PipelineStatus#CANCELING` that survives the same restore 
(`SubPlan#restorePipelineState`'s `CANCELING` handling in `stateProcess()` 
simply re-sends `task.cancel()` to every task). Savepoint has no pipeline-level 
equivalent of `CANCELING` to survive on — the only markers that a savepoint was 
in progress were the job-level `DOING_SAVEPOINT` status and the transient, 
in-memory `PendingCheckpoint`, and both are discarded unconditionally by the 
restore path above.
   
   ### How was this patch tested?
   
   **New E2E test** 
(`CheckpointCoordinatorFailoverIT#testStreamJobResolvesToSavepointDoneAfterMasterFailoverDuringSavepointShutdown`):
 submits a STREAMING job, waits for it to be `RUNNING` with output flowing, 
triggers a client savepoint asynchronously, then white-box polls 
`CheckpointCoordinator`'s private `pendingCheckpoints` field (same reflection 
idiom this module's own `CheckpointCoordinatorTest` unit tests already use, via 
`ReflectionUtils`) until the job's status has flipped to `DOING_SAVEPOINT` 
**and** the coordinator holds a savepoint checkpoint that has been dispatched 
but not yet fully acknowledged — i.e. genuinely mid-shutdown, neither 
un-started nor already complete. It kills the master at exactly that instant 
(mirroring the master-kill technique of the two existing tests in this file) 
and asserts the job resolves to `SAVEPOINT_DONE` within a bounded wait.
   
   Before the fix, this test reliably (not flakily) reproduced the job getting 
stuck at `RUNNING` — confirmed by running it against the unpatched restore path 
first.
   
   **Fix** (mirroring #10836's own persist-then-redrive pattern):
   - `JobMaster` gains a `redriveSavepointAfterRestore` flag, set by 
`CoordinatorService` when a restored job's prior status was `DOING_SAVEPOINT`.
   - `CoordinatorService#redriveSavepointAfterMasterFailover` consumes that 
flag concurrently with `JobMaster#run()` (which blocks until the job ends, so 
scheduling this after `run()` returns would never run while the job is active) 
and retries `JobMaster#savePoint()` until it succeeds, the job reaches a 
terminal state some other way, or a 5-minute budget elapses.
   - A single, unretried attempt gated only on the job's top-level status 
reaching `RUNNING` is not sufficient: `CheckpointCoordinator`'s own 
`isAllTaskReady` flag can still be `false` at that point (reset by 
`restoreCoordinator` whenever a task needed redeployment) until every task 
independently re-reports `READY_START` — which happens strictly after the 
coarser job-level status already cascaded to `RUNNING`. This was also confirmed 
empirically: a first version of the fix without retry still failed 
intermittently, hitting `TASK_NOT_ALL_READY_WHEN_SAVEPOINT` on its single 
attempt. `JobMaster#savePoint()` already recovers from that precondition 
without failing the job, so retrying is safe and never leaves the job worse off 
than before this fix.
   
   **Verification evidence:**
   - New test run in isolation: passed repeatedly (20+ runs) after the fix.
   - Both pre-existing tests in the same file 
(`testBatchJobCompletesAfterMasterFailover`, 
`testStreamJobContinuesAfterMasterFailover`) run together with the new one: 
`Tests run: 3, Failures: 0, Errors: 0` — no regression to the scenarios #10836 
already covers.
   - `./mvnw spotless:apply` run on both changed modules 
(`seatunnel-engine-server`, `connector-seatunnel-e2e-base`) before committing.
   
   A rarer (~10% in repeated stress-runs, not reproduced after system load 
dropped), apparently pre-existing race was also observed during this 
investigation: `runningJobStateIMap.get(jobId)` can occasionally read back as 
effectively bypassed (the job already present in `runningJobMasterMap` before 
`restoreJobFromMasterActiveSwitch` would run for it), skipping the 
restore-and-redrive path for that attempt. This reproduced identically whether 
the new test's fix was present or not (i.e. it is not something this PR 
introduces), affects the general master-switch restore reconciliation rather 
than anything savepoint-specific, and was not chased further to keep this PR's 
scope to the confirmed, well-understood savepoint gap. Flagging it here for 
awareness in case it is worth a dedicated investigation.
   
   ### Does this PR introduce _any_ user-facing change?
   
   No. This restores the originally-intended behavior of a client-triggered 
savepoint surviving a master failover during its shutdown phase.
   
   ### Check list
   
   * [ ] If any new Jar binary package adding in your PR, please add License 
Notice
   * [ ] If necessary, please update the documentation
   * [ ] If necessary, please update `incompatible-changes.md`
   
   🤖 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