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]
