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

   ## What this guards
   
   Two related fixes both concern job restore getting stuck after a Zeta master 
switch:
   
   - **#10692** ("[Fix][Zeta] Prevent terminal-state zombie jobs from being 
restored after master switch"): before this fix, 
`CoordinatorService#restoreAllRunningJobFromMasterNodeSwitch` funneled *every* 
entry found in `runningJobInfoIMap` — including jobs that had already reached a 
terminal status such as FINISHED/CANCELED — through the same `while 
(getResourceManager().workerCount(...) == 0)` wait loop that live 
(still-running) jobs legitimately need. A terminal job's IMap tombstone cleanup 
does not need a worker at all, so it could be starved indefinitely behind that 
loop. The fix added a pre-filter that resolves terminal-state entries 
immediately, before the worker-wait loop is ever reached, so only genuinely 
live jobs are subject to it.
   - **#10562** and **#10842**: a related pair of fixes in the same restore 
path — `RetryableHazelcastException` (thrown when an IMap partition is still 
loading, e.g. during a master switch) wasn't classified as retryable, so 
restore could fail outright instead of retrying; and a call site re-read 
`runningJobInfoIMap.get(jobId)` redundantly instead of reusing the 
already-available `jobInfo` parameter.
   
   I traced current `dev` HEAD 
(`CoordinatorService.restoreAllRunningJobFromMasterNodeSwitch`, 
`restoreJobFromMasterActiveSwitch`) to confirm what's still in place before 
writing this test:
   
   - The **terminal-job pre-filter from #10692 is present and unchanged in 
shape**: it iterates the restore candidates, and for any whose last known 
`JobStatus.isEndState() == true`, calls `restoreJobFromMasterActiveSwitch` 
immediately and removes it from the list — *before* the `while 
(getResourceManager().workerCount(Collections.emptyMap()) == 0)` loop that only 
the remaining (live) jobs fall through to.
   - The **`RetryUtils.retryWithException(...)` wrapping from #10562/#10842 
around the IMap reads is present** in both 
`restoreAllRunningJobFromMasterNodeSwitch` (fetching the candidate list) and 
`restoreJobFromMasterActiveSwitch` (fetching job state), and 
`restoreJobFromMasterActiveSwitch` correctly reuses the `jobInfo` parameter 
rather than re-reading the IMap.
   - The **worker-wait loop itself is still unbounded** — no timeout, no max 
iteration count:
     ```java
     // waiting have worker registered
     while (getResourceManager().workerCount(Collections.emptyMap()) == 0) {
         try {
             logger.info("Waiting for worker registered");
             Thread.sleep(1000);
         } catch (InterruptedException e) { ... }
     }
     ```
     This is a **related but distinct, still-open architectural gap**: if no 
worker ever registers after a master switch, restore of every remaining *live* 
job blocks forever (not just the terminal ones #10692 already excludes). This 
PR is test-only and does not fix that gap — it only proves that live jobs 
correctly wait (and successfully proceed once a worker appears), while terminal 
jobs correctly do *not* wait at all. Flagging it here per the task brief as a 
separate, out-of-scope observation.
   
   Existing coverage for #10692 was 
`CoordinatorServiceTest#testTerminalZombieJobShouldNotRestartAfterMasterSwitch`,
 a single-job scenario that hand-injects stale terminal IMap entries via direct 
`IMap.put(...)` calls and reflection into a 2-node unit-test-style cluster. It 
does not exercise a **mix** of terminal and live jobs restoring together, and 
does not construct the "zero workers registered at switch time" precondition 
explicitly. This PR adds a real multi-node E2E path for that gap.
   
   ## What the new test does
   
   
`SplitClusterPendingJobLifecycleFailoverIT#testTerminalJobCleanupSkipsWorkerWaitAfterMasterSwitch`:
   
   1. Starts two master-eligible nodes plus one temporary worker; runs a small 
batch job (`cluster_batch_fake_to_localfile_template.conf`) to completion on 
that worker, so it reaches FINISHED the normal way (leaving a real 
`pendingJobCleanupIMap` tombstone behind, not a hand-injected one).
   2. Shuts the worker down — **only after** that batch job's tasks already 
reached a terminal execution state — then confirms (via cluster membership 
size) that zero workers remain anywhere in the cluster.
   3. Submits a second job with zero workers registered. 
`ScheduleStrategy#WAIT` resource pre-checks fail gracefully rather than 
erroring out (see `JobMaster#preApplyResources`), so this job lands in PENDING: 
a genuinely live, non-end-state job that legitimately still needs a worker.
   4. Shuts down the active master, confirming the standby (which has held a 
live IMap backup since before either job was submitted) becomes the new active 
coordinator with zero registered workers.
   5. Asserts, for a 10-second stability window while zero workers remain 
registered: the terminal job's status stays FINISHED (never re-queued), while 
the live job is **not yet** in the new master's pending-job queue — proving 
they take genuinely different code paths, not just "both happened to resolve 
quickly."
   6. Starts a fresh worker; asserts the live job's restore now proceeds and 
reaches RUNNING, and re-confirms the terminal job's status is unaffected.
   
   ## Why the trigger construction is reliable
   
   The precondition ("terminal job present, live job present, and the *new* 
active master has zero registered workers") is built purely by controlling node 
start/stop **order**, not by racing or sleeping past a hoped-for window:
   
   - Both master nodes are started before either job is submitted, so the IMap 
backup that survives the switch is already live.
   - The temporary worker is torn down **only after** the batch job's tasks 
have already completed — not while any task is still 
DEPLOYING/RUNNING/CANCELING on it. This matters because 
`CoordinatorService#failedTaskOnMemberRemoved` is a *synchronous* Hazelcast 
membership listener: killing a worker out from under an in-flight task fails 
that task immediately (not after a delay), and with the default 
`job.retry.times`/`job.retry.interval.seconds` (3/3) a job that loses its only 
worker mid-flight would flip to terminal FAILED roughly 9 seconds later. That 
would have silently collapsed this test's mixed-state precondition (both "jobs" 
ending up terminal) — I verified this failure path directly in 
`CoordinatorService` and had it independently confirmed by another 
investigation before choosing this construction.
   - The "live" job is therefore built as a genuinely non-end-state PENDING job 
submitted *after* the sole worker is already gone, rather than a RUNNING job 
whose worker is killed later — this exercises the exact same `isEndState() == 
false` branch in `restoreAllRunningJobFromMasterNodeSwitch` (the fix in #10692 
doesn't distinguish RUNNING from PENDING; it only checks end-state-ness), 
without any dependency on winning a race against `failedTaskOnMemberRemoved`.
   - No sleeps are used to "hope" the precondition holds; every gating 
condition (cluster member counts, job statuses) is asserted via `Awaitility` 
before the next step proceeds.
   
   ## Why scenario 2 (transient IMap-loading exception classification) is not 
covered here
   
   I considered also covering the #10562/#10842 
retry-on-`RetryableHazelcastException` behavior, but this test module's 
Hazelcast IMaps have no `MapStore` configured (no external/S3-backed loading, 
confirmed by inspecting the test `hazelcast.yaml`/`seatunnel.yaml` resources), 
so there is no honest, black-box way to make a real IMap partition report 
"still loading" during a master switch in this environment. The only existing 
coverage of that exact mechanism 
(`CoordinatorServiceTest#testRestoreUsesProvidedJobInfoInitializationTimestamp`)
 already does this correctly via a Mockito spy that throws 
`RetryableHazelcastException` from an IMap `get(...)` call — that is an 
appropriate way to test an internal exception-handling contract, but not 
something a black-box multi-node E2E test can reproduce honestly without 
equivalent internal instrumentation. Per the task guidance, I focused this PR 
on the scenario I could construct reliably (terminal-vs-live job restore 
gating) rather than force 
 a weaker or flaky simulation of the IMap-loading race.
   
   ## Test plan
   
   - `./mvnw spotless:apply -pl 
seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base -am -nsu` — 
BUILD SUCCESS, no formatting changes needed beyond what's in this diff.
   - `./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; directly confirmed (via `javap`) that the new test 
method and its helper compiled into 
`SplitClusterPendingJobLifecycleFailoverIT.class`, since `-DskipTests` alone 
(Surefire-only) still compiles test sources, unlike `-Dmaven.test.skip=true` 
which would have skipped compilation entirely.
   - Actual test execution is left to CI, consistent with this repository's 
local-verification policy for the main Apache SeaTunnel repo.
   


-- 
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