goutamadwant opened a new issue, #12010:
URL: https://github.com/apache/seatunnel/issues/12010

   ### Search before asking
   
   - [x] I searched existing issues and pull requests and found no report 
covering this submission-acknowledgement case.
   
   Related work such as #7693, #9532, and #11653 covers pending-job scheduling 
and failover, but not an interrupted queue insertion being acknowledged as a 
successful submission.
   
   ### What happened
   
   On the current `dev` branch, `CoordinatorService.submitJob` initializes the 
`JobMaster`, creates a `PendingJobInfo`, calls `pendingJobQueue.put(...)`, and 
then completes `jobSubmitFuture` successfully.
   
   `PeekBlockingQueue.put` catches `InterruptedException`, logs it, and returns 
`void`. The interrupt is not restored and the failure is not propagated to 
`CoordinatorService`. As a result, the submission path continues as if 
insertion succeeded:
   
   1. `jobSubmitFuture.complete(null)` reports success to the caller.
   2. The physical plan is updated to `PENDING`.
   3. The coordinator logs that the job entered the pending queue.
   4. The pending queue does not contain the job.
   
   This can occur during a coordinator shutdown or master transition. 
`clearCoordinatorService()` calls `executorService.shutdownNow()`, which can 
interrupt an in-flight submission worker immediately before the queue 
insertion. Although the backing `LinkedBlockingQueue` is unbounded, its `put` 
operation acquires an internal lock interruptibly and can therefore throw 
before inserting.
   
   Distributed job metadata is written before the insertion, so a later active 
master may restore the job. This report does not assume permanent job loss. The 
confirmed problem is that the original submission can be reported as successful 
even though the active coordinator did not accept the job into its pending 
queue, making execution dependent on later recovery.
   
   Relevant code:
   
   - 
[`PeekBlockingQueue.put`](https://github.com/apache/seatunnel/blob/b03cf1f1697ed84d0f57f416673cec40a075ef5a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/utils/PeekBlockingQueue.java#L59-L70)
   - 
[`CoordinatorService.submitJob`](https://github.com/apache/seatunnel/blob/b03cf1f1697ed84d0f57f416673cec40a075ef5a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java#L1441-L1463)
   - 
[`CoordinatorService.clearCoordinatorService`](https://github.com/apache/seatunnel/blob/b03cf1f1697ed84d0f57f416673cec40a075ef5a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java#L1280-L1330)
   
   ### Expected behavior and proposed contract
   
   A successful submission response should mean that the job was accepted into 
the pending scheduling path, or into another explicitly guaranteed recovery 
path.
   
   Before preparing a fix, please confirm which contract is intended:
   
   1. **Fail interrupted or stale submission:** if interruption occurs before 
the job reaches an accepted scheduling or recovery boundary, complete 
`jobSubmitFuture` exceptionally with a retriable failure and preserve the 
thread interrupt status. Cleanup must interrupt the initialized `JobMaster` and 
remove partial state only when it is still owned by this submission, without 
deleting state already claimed or restored by a new active master.
   2. **Make acceptance non-interruptible and lifecycle-safe:** because the 
queue is unbounded, use a non-interruptible insertion operation, but coordinate 
insertion with the coordinator activation/shutdown epoch so no job can be 
enqueued after `clearCoordinatorService()` has cleared the local queue. 
Acknowledge success only after the queue and job-id index are updated 
consistently, or after another explicitly documented durable-recovery boundary 
has been reached.
   
   `restoreJobFromMasterActiveSwitch()` also uses `pendingJobQueue.put(...)`. 
If the queue contract changes to propagate interruption, the restore path needs 
an explicit interruption and epoch policy as well.
   
   ### Reproduction
   
   The behavior is deterministic when the submission task is paused immediately 
before the real queue insertion and `clearCoordinatorService()` interrupts the 
coordinator executor:
   
   The injected queue deliberately widens the narrow production race window 
immediately before `PeekBlockingQueue.put` acquires its outer lock. It still 
uses the real coordinator executor and the production 
`clearCoordinatorService()` shutdown path. The interrupt comes from 
`executorService.shutdownNow()`, not from the test thread.
   
   ```java
   private static class InterruptOnShutdownPendingJobQueue
           extends PeekBlockingQueue<PendingJobInfo> {
       private final CountDownLatch putStarted = new CountDownLatch(1);
       private final CountDownLatch neverReleased = new CountDownLatch(1);
   
       private InterruptOnShutdownPendingJobQueue() {
           super(PendingJobInfo::getJobId);
       }
   
       @Override
       public void put(PendingJobInfo element) {
           putStarted.countDown();
           try {
               neverReleased.await();
           } catch (InterruptedException e) {
               Thread.currentThread().interrupt();
           }
           super.put(element);
       }
   }
   
   InterruptOnShutdownPendingJobQueue queue = new 
InterruptOnShutdownPendingJobQueue();
   ReflectionUtils.setField(coordinatorService, "pendingJobQueue", queue);
   
   PassiveCompletableFuture<Void> submitFuture =
           coordinatorService.submitJob(jobId, data, false);
   assertTrue(queue.putStarted.await(20, TimeUnit.SECONDS));
   
   coordinatorService.clearCoordinatorService();
   
   assertDoesNotThrow(submitFuture::join);
   assertFalse(queue.contains(jobId));
   assertEquals(0, queue.size());
   assertTrue(runningJobInfoIMap.containsKey(jobId));
   assertEquals(JobStatus.PENDING, runningJobStateIMap.get(jobId));
   ```
   
   The focused `CoordinatorServiceTest` passed while asserting this 
inconsistent result on both JDK 8 and JDK 11:
   
   ```shell
   ./mvnw -pl seatunnel-engine/seatunnel-engine-server \
     -Dskip.spotless=true -Dcheckstyle.skip=true -Dlicense.skip=true \
     
-Dtest=CoordinatorServiceTest#reproduceInterruptedPendingJobInsertionReportsSuccess
 test
   ```
   
   ### Observed log
   
   ```log
   ERROR PeekBlockingQueue - Put element into queue failed. 
java.lang.InterruptedException
   INFO  PhysicalPlan - Job interrupted_pending_job_insertion (...) turned from 
state CREATED to PENDING.
   INFO  CoordinatorService - The submit job enter the pending queue ...
   ```
   
   At that point, the submission future has completed normally while the 
pending queue has no entry for the job. The distributed `runningJobInfoIMap` 
entry still exists and `runningJobStateIMap` reports `PENDING`, which is why 
later coordinator recovery may be possible even though the original acceptance 
response was false.
   
   ### SeaTunnel Version
   
   Current `dev` at `b03cf1f1697ed84d0f57f416673cec40a075ef5a`.
   
   ### SeaTunnel Config
   
   ```conf
   The issue is not configuration-specific. The reproduction uses the existing
   seatunnel-engine-server test resource: batch_fake_to_console.conf.
   ```
   
   ### Running Command
   
   ```shell
   ./mvnw -pl seatunnel-engine/seatunnel-engine-server \
     -Dskip.spotless=true -Dcheckstyle.skip=true -Dlicense.skip=true \
     
-Dtest=CoordinatorServiceTest#reproduceInterruptedPendingJobInsertionReportsSuccess
 test
   ```
   
   ### Error Exception
   
   ```log
   java.lang.InterruptedException
       at 
java.util.concurrent.locks.AbstractQueuedSynchronizer.acquireInterruptibly(...)
       at java.util.concurrent.LinkedBlockingQueue.put(...)
       at 
org.apache.seatunnel.engine.server.utils.PeekBlockingQueue.put(PeekBlockingQueue.java:62)
       at 
org.apache.seatunnel.engine.server.CoordinatorService.lambda$submitJob$11(CoordinatorService.java:1450)
   ```
   
   ### Zeta or Flink or Spark Version
   
   Zeta engine from the current `dev` branch.
   
   ### Java or Scala Version
   
   Reproduced independently with JDK 8 and JDK 11.
   
   ### Are you willing to submit a PR?
   
   Yes, after maintainers confirm the intended interruption contract.
   
   ### Code of Conduct
   
   - [x] I agree to follow the project's Code of Conduct.
   


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