davidzollo opened a new pull request, #12035:
URL: https://github.com/apache/seatunnel/pull/12035
## What this guards
This adds an E2E regression test for the fix in
`e676f46dc73a49e410661ccb49b325d6871eb70a` ("[Fix][Zeta] Fix
stop-with-savepoint hang in DOING_SAVEPOINT", #11489).
**The bug**: `JobMaster#savePoint()` moved the job into `DOING_SAVEPOINT`
and then triggered every pipeline's `CheckpointCoordinator#startSavepoint()`.
If a pipeline's coordinator rejected the request's precondition
(`TASK_NOT_ALL_READY_WHEN_SAVEPOINT` when not all subtasks had reported
`READY_START` yet, or `CHECKPOINT_COORDINATOR_SHUTDOWN` if the coordinator was
mid-shutdown), nothing ever rolled the job back to `RUNNING` or drove it to a
terminal state — the job hung in `DOING_SAVEPOINT` forever.
**The fix** classifies every pipeline's savepoint-completion outcome in
`JobMaster#waitSavepointCompleted`/`isSavepointStartPreconditionFailure`. When
every failure is a pre-start precondition rejection and no pipeline reached
`SUSPEND`, `restoreRunningAfterSavepointStartFailure()` restores the job to
`RUNNING` so the caller can retry. Any other failure instead drives the job to
a terminal state through the new `PhysicalPlan#savepointFailed()` fallback.
## What coverage was missing
The fix's own accompanying tests (`JobMasterTest`,
`SavePointBusySourceTest`) are unit-level: they invoke the classification
methods directly via reflection against fabricated, already-failed
`CompletableFuture`s — they never submit a real job or drive a real
`CheckpointCoordinator` through its actual precondition check.
`SavepointRestoreIT` and `CheckpointRestoreWithStopIT` in the same package only
exercise the happy path (savepoint succeeds, then restore from it). None of the
existing coverage proves that a real, client-triggered savepoint request
landing on a real `CheckpointCoordinator` while it genuinely rejects the
not-ready precondition actually recovers instead of hanging.
## What the new test does
`SavepointPreconditionRecoveryIT` (in-JVM split master/worker cluster,
following the house style of
`SplitClusterPendingJobLifecycleFailoverIT`/`CheckpointCoordinatorFailoverIT`):
1. Submits a real streaming FakeSource → LocalFile job via `SeaTunnelClient`.
2. Spin-polls the master's in-process `CheckpointCoordinator` (reached
through `SeaTunnelServer`/`CoordinatorService`, exactly as other tests in this
package already do) via reflection on its `isAllTaskReady` field — the exact
flag `startSavepoint()` checks — until the coordinator exists but has not yet
observed every subtask report `READY_START`.
3. Fires a real `engineClient.savePointJob(jobId)` in the very next
statement so it arrives on the master while the precondition is still false.
4. Asserts the rejection is specifically the
`TASK_NOT_ALL_READY_WHEN_SAVEPOINT` precondition (proving the intended race was
actually hit, not some unrelated failure).
5. **Core regression assertion**: asserts the job recovers to `RUNNING`
within a bounded timeout instead of hanging in `DOING_SAVEPOINT`.
6. Asserts the pipeline keeps producing data after recovery (proving the
recovery is functionally real, not a cosmetic status flip).
7. Asserts a second, unraced savepoint (issued once `isAllTaskReady` has
settled true) now succeeds normally and drives the job to `SAVEPOINT_DONE` —
proving the checkpoint coordinator itself is fully healthy after the earlier
rejection, not just superficially reporting `RUNNING`.
## Why the trigger is reliable
`CheckpointCoordinator` objects are created synchronously inside
`JobMaster#init()`, and by the time they exist the job is already registered in
`CoordinatorService#runningJobMasterMap` (that registration strictly
happens-before the asynchronous `JobMaster#run()`/`init()` call that creates
the coordinators). So a savepoint request fired the moment the coordinator
becomes visible is guaranteed to reach the real precondition check in
`CheckpointCoordinator#startSavepoint()` — not rejected earlier for the
unrelated "job not running" reason. Task deployment, worker startup, and the
`READY_START` report round trip all happen strictly after coordinator creation
and take real wall-clock time (worker RPC, classloading, thread startup),
giving the white-box poll a wide safety margin over blind timing.
## Test plan
- `./mvnw spotless:apply -pl
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` —
passed.
- `./mvnw install -pl
'seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base,!seatunnel-engine/seatunnel-engine-ui'
-am -nsu -Dmaven.gitcommitid.skip=true -DskipTests -Dspotless.check.skip=true
-T 3C` — `BUILD SUCCESS`, with
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base` module
reporting `SUCCESS`. Directly confirmed the new test class actually compiled:
`target/test-classes/org/apache/seatunnel/engine/e2e/SavepointPreconditionRecoveryIT.class`
exists. (`seatunnel-engine-ui` excluded since it is unmodified by this change
and pulling it into the reactor triggers an unrelated `npm install` that has no
network access in this environment; its already-published `.m2` artifact is
used instead.)
- Test execution itself was not run locally per this task's
local-verification policy for the Apache SeaTunnel repository; correctness of
the race/assertions was verified by tracing the exact source paths
(`CheckpointCoordinator#startSavepoint`,
`JobMaster#waitSavepointCompleted`/`isSavepointStartPreconditionFailure`/`restoreRunningAfterSavepointStartFailure`,
`CoordinatorService#savePoint`/`pendingJobSchedule`) this test exercises. CI
on this PR is the source of truth for actual execution.
🤖 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]