DanielLeens commented on PR #11757:
URL: https://github.com/apache/seatunnel/pull/11757#issuecomment-5605347445

   ## Pushed `c337f91`: generation-tagged async/timer registrations, closing F2 
(plus F1/F4, F5, F7/F8)
   
   @SEZ9 @waterWang this is the commit that addresses the open blocker. Diff is 
still limited to `TaskExecutionService.java` and 
`TaskExecutionServiceTest.java`.
   
   **F2 (the blocker) - stale branch now cancels exactly its own futures.** 
`taskAsyncFunctionFuture` and `timerFlushFutures` stay keyed by 
`TaskGroupLocation`, but every entry is now an `OwnedFuture` that records the 
`TaskGroupContext` that was active for the location at registration time 
(`asyncExecuteFunction()` / `registerTimerFlushTask()` read 
`executionContexts.get(location)`; the owner is `null` when no generation is 
active, e.g. a barrier arriving for an already-finished group). Teardown goes 
through one helper, `detachAsyncResources(location, ownedContext, 
ownsLocation)`:
   - a tracker that still owns the location detaches everything under it (own 
entries, untagged entries, and leftovers of superseded generations that never 
reached their own teardown) - identical to the pre-PR sweep for the normal path;
   - a superseded tracker detaches only the entries tagged with its own 
`ownedContext` and leaves the active generation's entries in place.
   
   `finishOwnedResources()` uses it in both branches (the stale branch no 
longer returns early after `recycleClassLoader`), and `cancelAllTask()` uses 
the same rule instead of skipping the cleanup for a superseded generation, so 
both the completionLatch==0 path and the FAILED-fallthrough / 
cancellationFuture path are covered.
   
   **F1/F4 - nothing heavy under the service monitor any more.** The critical 
section under `TaskExecutionService.this` now only does the ownership check and 
the map mutations (`finishExecutionContext`, `cancellationFutures.remove`, 
detach). `recycleClassLoader()` and the actual `Future#cancel` calls run after 
the monitor is released, on state that is no longer reachable through the maps. 
Async functions keep `cancel(true)`, timer flushes keep `cancel(false)`.
   
   **F7/F8 - invariant documented and enforced, warning made accurate.** The 
`executionContexts` Javadoc states that every writer (`deployLocalTask()` 
publish, `finishExecutionContext()` removal) must hold the service monitor; 
`finishExecutionContext()` and `detachAsyncResources()` call 
`ensureServiceMonitorHeld()` (`Thread.holdsLock`, same pattern as 
`CheckpointCoordinator`). With the invariant enforced, 
`finishExecutionContext()` does the identity check followed by a plain 
`remove(key)` instead of the equals()-based `remove(key, value)`, and the 
Javadoc keeps the explanation of why the Lombok `@Data` equals must not be the 
ownership test. The stale-branch log now says whether the location has no 
active context or is owned by a newer execution context, instead of always 
claiming a newer generation exists.
   
   **F6** - added a comment at the publish block explaining why the 
`cancellationFutures.put` cannot displace a live future when reached through 
`deployTask()`; no code change, per the earlier trace. **F3** stays with #12164 
/ #12218.
   
   **F5 / tests.** `testStaleTaskDoneReleasesOnlyOwnGenerationResources` and 
`testStaleFailedTaskDoneCancelsOnlyOwnGenerationResources` (renamed from the 
two `...DoesNotCleanupNewerGenerationResources` tests) now register real async 
functions and timer flushes for both generations through the public API - old 
generation active first, then the new generation takes over the same location - 
and assert that the old generation's futures are cancelled while the new 
generation's context, cancellation future, async function and timer stay live. 
`testActiveTaskDoneSweepsAllLocationResources` covers the owning branch, 
including an untagged registration made before any generation owned the 
location. All three remove their injected entries from the shared service in a 
`finally` block.
   
   **Verification.** Local: `spotless:apply` + `spotless:check` on 
`seatunnel-engine-server` only, per the usual rule for this repo (no local 
build or test run); compile, unit tests and E2E are left to the CI run on 
`c337f91`, which is the only verification that counts here. Heads-up for 
whoever merges second: #12218 (F3 fix) touches the same `deployLocalTask()` / 
`taskDone()` regions, so one of the two will need a rebase.
   


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