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

   ### Purpose of this pull request
   
   Fixes the hang in #11679, where a job cancelled by a downstream privilege 
check left the
   `seatunnel.sh` client blocked indefinitely and burned a full 180-minute CI 
slot.
   
   I reproduced it in CI and took thread dumps of both the engine and the 
blocked client while it was
   stuck. Dumps and the coordinator log are attached to the issue. The client 
turned out to be a
   symptom — it never receives a terminal result because the engine deadlocks 
before it can publish
   one.
   
   **Root cause.** `TaskGroupLocation` is `{jobId, pipelineId, taskGroupId}` 
and is reused verbatim
   across pipeline restore generations. `TaskGroupExecutionTracker.taskDone()` 
removes that location's
   entry from `executionContexts`, so a late `taskDone()` from an earlier 
generation can delete the
   context that the current generation's `deployLocalTask()` has just installed.
   
   `BlockingWorker.run()` then resolved its class loader through 
`executionContexts.get(location)`
   *before* its `try` block, while `startedLatch.countDown()` sat *inside* it:
   
   ```java
   ClassLoader classLoader =
           executionContexts
                   
.get(taskGroupExecutionTracker.taskGroup.getTaskGroupLocation())  // may be null
                   .getClassLoaders()                                           
      // NPE here
                   .get(tracker.task.getTaskID());
   ...
   try {
       startedLatch.countDown();   // never reached
   ```
   
   A missing context therefore threw before the latch was ever counted down. 
The exception was
   swallowed by the submitting `Future`, so nothing was logged, and 
`submitBlockingTask()` waited on
   `startedLatch.await()` forever — **while holding the `SubPlan` monitor**. 
The checkpoint-error
   thread then blocked at the `synchronized updatePipelineState()` entry and 
could never move the
   pipeline to a terminal state, so `JobMaster`'s completion future never 
completed and the client
   waited on it indefinitely.
   
   From the dumps:
   
   ```
   "seatunnel-coordinator-service-8"  WAITING (parking)
       at java.util.concurrent.CountDownLatch.await(CountDownLatch.java:231)
       at 
...TaskExecutionService.submitBlockingTask(TaskExecutionService.java:405)
       ...
       at ...SubPlan.updatePipelineState(SubPlan.java:397)
       - locked <0x00000000f6d2c470> (a SubPlan)
   
   "hz.main.generic-operation.thread-96"  BLOCKED (on object monitor)
       at ...SubPlan.updatePipelineState(SubPlan.java:345)
       - waiting to lock <0x00000000f6d2c470> (a SubPlan)
       at ...SubPlan.handleCheckpointError(SubPlan.java:644)
   ```
   
   ### Does this PR introduce a user-facing change?
   
   No behaviour change in the normal path. A deployment whose execution context 
has disappeared now
   fails with a clear `IllegalStateException` naming the task group, instead of 
hanging silently.
   
   ### How was this patch tested?
   
   Added `TaskDeployStaleContextRaceTest`, which asserts the narrow contract 
that prevents the hang:
   deploying a task group must return even if its execution context disappears 
while the deployment is
   in flight. It does not assert *how* it returns — success or failure both 
satisfy it. Only a
   deployment that never returns is the defect.
   
   Run on the JDK 8 leg, same commit, test alone versus test plus fix:
   
   | | Result |
   |---|---|
   | Test on unpatched engine | `Tests run: 1, Failures: 1` — *"deployTask did 
not return within 30s on iteration 0"*, 32.8s |
   | Test with this fix | `Tests run: 1, Failures: 0` — 30 iterations in 
**3.9s** |
   
   The original `PaimonWithS3IT` hang reproduced roughly 1 in 6 runs before 
this change.
   
   ### A note on the fix
   
   `startedLatch.countDown()` deliberately stays **before** `init()`. Moving it 
into `finally`
   unguarded would make `submitBlockingTask()` wait for tasks to *finish* 
rather than *start*, which
   would deadlock every long-running task. The signal is guarded by a 
`startSignalled` flag so each
   worker releases the latch exactly once, including on the failure path.
   
   ### Scope, and a follow-up
   
   This fixes the **hang**, not the race that triggers it. A stale `taskDone()` 
can still remove a
   live execution context; the deployment now fails loudly instead of wedging 
the engine. Making that
   removal conditional — so a tracker from an earlier generation cannot delete 
a newer entry — is a
   separate change with its own risk profile, and I would rather propose it on 
its own than bundle it
   here. Happy to open it as a follow-up if maintainers agree with the 
direction.
   
   Marked as draft for that reason: the containment is verified, the locking 
design decision is yours.
   
   Related: #11718 keeps the CI blast radius small and is complementary — it 
does not overlap with this
   change.
   
   ### Check list
   
   * [x] Code changed are covered with tests, or it does not need tests for 
reason
   * [x] If any new Jar binary package adding in your PR, please add License 
Notice according
     [New License 
Guide](https://github.com/apache/seatunnel/blob/dev/docs/en/contribution/new-license.md)
   * [x] If necessary, please update the documentation to describe the new 
feature.
   * [x] If you are contributing the connector code, please check that the 
following files are updated:
   * [x] Update the `docs/en/connector-v2` and `docs/zh/connector-v2` for the 
connector.
   


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