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]