Vivek1106-04 commented on PR #12511:
URL: https://github.com/apache/seatunnel/pull/12511#issuecomment-5982000557

   @DanielLeens thanks for the thorough review. All eight points are addressed 
in `9143d479c` and `a691fb195`. The branch is also rebased onto `dev` 
`af80a6704`, so it is 0 behind and picks up #12444 for the PayPal flake.
   
   **Issue 1 (gate has no test that fails without it).** New 
`testTriggerAfterInFlightCheckpointEndsButBeforeSavepointIsCreatedRearms` uses 
your lock-holding idea, with a real drain thread. The test thread takes the 
coordinator `lock` and drops `pendingCounter` to 0. It then waits until the 
savepoint caller has woken and is `BLOCKED` on the lock, and fires a trigger of 
each type; the lock is reentrant, so those calls run. Each trigger must create 
nothing and re-arm. After the lock is released, exactly one savepoint is 
created. With `|| savepointDraining` removed it fails: `a CHECKPOINT_TYPE 
trigger created a checkpoint ahead of the savepoint ==> expected: <0> but was: 
<1>`. @Rangsh ported the matching #12493 test, thanks. I went with this version 
instead because it covers the same window without setting the gate by 
reflection, so no gate is left without a request behind it.
   
   **Issue 2 (re-interrupt on a pool worker).** The flag is now re-set only 
when the caller is not a `ForkJoinWorkerThread`. 
`testInterruptedSavepointDrainFailsTheRequestAndResumesTriggering` asserts that 
a caller on its own thread keeps the flag. New 
`testInterruptedSavepointDrainOnPoolWorkerClearsTheInterruptFlag` runs the 
drain on a `ForkJoinPool` and asserts that the flag is clear when 
`startSavepoint()` returns. It fails with the unconditional interrupt 
(`expected: <false> but was: <true>`).
   
   **Issue 3 (cancel during the restore wait ends FAILED).** Your inference was 
right, and I confirmed it by running it. New 
`CheckpointErrorRestoreEndTest#testCancelDuringPipelineRestoreWaitEndsTheJobCanceled`
 fails the first checkpoint, waits until the pipeline is reset for restore (30 
s wait in its conf), and cancels:
   - previous head `91001c42c`: `expected: <CANCELED> but was: <FAILED>`
   - `dev`: passes. The log shows `PhysicalPlan` going `RUNNING -> CANCELING`, 
then `cancelPipeline` waiting on the SubPlan monitor until the wait ends. The 
pipeline restarts (`CREATED -> SCHEDULED -> DEPLOYING -> RUNNING`), is 
cancelled straight away, and the job ends `CANCELED`.
   
   I kept `dev`'s outcome. If the job is `CANCELING` when the restore is 
abandoned, the pipeline ends `CANCELED`; otherwise it ends in the state it 
failed with and keeps its error. The savepoint case is unchanged, and the 
cancelled job no longer restarts a pipeline only to cancel it.
   
   **Issue 4 (root cause and point-in-time check).** You were right that 
`forceStop` is not a no-op. I took a log trace of the restore-wait 
`SavePointTest` on `dev`'s `SubPlan`. During the wait, `forceStop` moves the 
enumerator and committer `CREATED -> CANCELED`, and both log `state process is 
not start`. Their vertex state process was stopped, so `taskFuture` is never 
completed. Only the source task's future completes. The pipeline therefore 
never counts all tasks ended. 2 s later the restart goes `SCHEDULED -> 
DEPLOYING -> RUNNING` and deploys nothing, because `DEPLOYING` starts only 
`CREATED` tasks. The pipeline then sits in `RUNNING`.
   
   So the lost stop comes from `PhysicalVertex.forceStop` on a task whose state 
process is not running, not from callbacks queued behind the monitor. The 
monitor does explain why `cancelPipeline` waits for the restart in the Issue 3 
scenario. The Javadoc and the PR description now describe this mechanism.
   
   On the residual window: a stop decided after the re-check but before deploy 
hits the same `forceStop` gap. A second `isNeedRestore()` check in `SCHEDULED` 
would narrow the window, but it is still check-then-act against a stop path 
that takes no SubPlan lock, so it cannot close it. The real fix is for 
`forceStop` to complete the future of a task whose state process is not 
running. That changes a stop path shared by every job, so I would rather do it 
as a separate issue and PR than grow this one. Happy to open it if you agree.
   
   **Issue 5 (wall-clock SavePointTest).** The conf now pins 
`job.retry.interval.seconds = 3`. The fixed `Thread.sleep(2000L)` is replaced 
by an await on the job being `RUNNING` with every pipeline's coordinator 
reporting all tasks ready, so the savepoint cannot hit 
`TASK_NOT_ALL_READY_WHEN_SAVEPOINT`.
   
   **Issue 6 (futures completed under the lock, no log).** 
`failDrainingSavepoint` clears the gate under the lock and completes the 
request after leaving it. `cleanPendingCheckpoint` does the same: it captures 
the draining request under the lock and fails it after the synchronized block. 
That is safe for the reset window you pointed out, because 
`restoreCoordinator()` clears `shutdown` only after `cleanPendingCheckpoint` 
returns, and by then the request is done. Interrupted and cut-short drains now 
`LOG.warn` with job id, pipeline id and close reason.
   
   **Issue 7.** `drainAndTriggerSavepoint` has a Javadoc stating the loop 
invariant. `savepointPendingCheckpoint` is now `volatile`, and its comment says 
that only tests read it.
   
   **Issue 8.** #12493 is closed in favour of this PR. I would keep the 
`SubPlan` change here: the two commits stay separate, and without it any 
correct drain fix turns `SavePointTest` red.
   
   **Coverage gap from 2.2.** The reset test is now parameterized over 
`CHECKPOINT_COORDINATOR_RESET`, `PIPELINE_END`, 
`CHECKPOINT_COORDINATOR_COMPLETED` and `CHECKPOINT_INSIDE_ERROR`, and each must 
fail the request with `CHECKPOINT_COORDINATOR_SHUTDOWN`.
   
   Local results on JDK 17 at `a691fb195`: `CheckpointCoordinatorTest` 25/25, 
`SavePointTest` 6 run + 1 existing `@Disabled`, `CheckpointErrorRestoreEndTest` 
2/2. `spotless:apply` makes no changes and `-DskipTests verify` passes. The red 
evidence for each new test is in the updated PR description. CI is running on 
the new head.
   


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