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

   This is my own PR, so per protocol I'm posting this as a plain issue comment 
(GitHub blocks self-approval/self-request-changes via the review API) and 
stating the merge determination explicitly in the text.
   
   Status check before the substance: I re-fetched the PR live. The head is 
still `3851e0d9f2b204fa0cb7d030d9d64321dbccd936` (unchanged since my last round 
yesterday at 2026-08-21T14:13:03Z), and no one has posted a new commit, 
comment, or review since then. So there is nothing new to react to — but since 
this is a self-authored PR, protocol asks me to re-do a full, independent read 
of the current diff rather than just restate the prior conclusion from memory. 
I did that below, re-reading `CoordinatorService.java`, 
`JobHistoryService.java`, and the new test file directly at the current head, 
and I land in the same place as yesterday's round.
   
   # What Problem Does This PR Solve?
   - User pain point: `JobHistoryService` registers three cluster-wide 
Hazelcast `EntryExpiredListener`s (on `IMAP_FINISHED_JOB_STATE`, 
`IMAP_FINISHED_JOB_METRICS`, `IMAP_FINISHED_JOB_VERTEX_INFO`) in its 
constructor, and `CoordinatorService.checkNewActiveMaster()` -> 
`initCoordinatorService()` constructs a brand-new `JobHistoryService` every 
time this node becomes the active master (`CoordinatorService.java:515-529`). 
Before this PR nothing ever deregistered the previous instance's listeners, so 
every active -> inactive -> active master transition (failover, split-brain 
healing, rolling restart) leaked three listener registrations, kept the old 
`JobHistoryService` (and everything it transitively references) reachable from 
the IMap listener registry, and caused a real finished-job expiration event to 
fire once per leaked instance — duplicating `CleanLogOperation` RPC fan-out and 
cleanup counters.
   - Fix approach: capture the `UUID` returned by each `addEntryListener(...)` 
call as `final` fields, add a best-effort `JobHistoryService.close()` 
(`JobHistoryService.java:184-195`) that deregisters exactly those three via a 
`removeEntryListenerQuietly` helper (`:450-457`, catches and logs, never 
throws), and invoke `close()` from 
`CoordinatorService.clearCoordinatorService()` 
(`CoordinatorService.java:1211-1213`) — the single method already used both 
when a node leaves the active-master role and during full coordinator shutdown. 
`jobHistoryService` is deliberately left non-null so in-flight readers keep 
working after the role switch.
   - One-sentence summary: closes a listener/side-effect leak on Zeta 
master-role transitions by tracking and explicitly deregistering the three 
finished-job expiration listeners when the coordinator service is cleared.
   
   # 1. Code Change Review
   ## 1.1 Core Logic Analysis
   
   **Before** (`JobHistoryService.java` constructor):
   ```java
   this.finishedJobStateImap.addEntryListener(
           new FinishedJobExpiredListener<>(Constant.IMAP_FINISHED_JOB_STATE), 
true);
   this.finishedJobMetricsImap.addEntryListener(
           new 
FinishedJobExpiredListener<>(Constant.IMAP_FINISHED_JOB_METRICS), true);
   this.finishedJobDAGInfoImap.addEntryListener(
           new JobInfoExpiredListener(Constant.IMAP_FINISHED_JOB_VERTEX_INFO), 
true);
   ```
   Registration ids discarded, no removal path anywhere.
   
   **After** (`JobHistoryService.java:154-162`):
   ```java
   this.finishedJobStateListenerId =
           this.finishedJobStateImap.addEntryListener(
                   new 
FinishedJobExpiredListener<>(Constant.IMAP_FINISHED_JOB_STATE), true);
   this.finishedJobMetricsListenerId = ...
   this.finishedJobDAGInfoListenerId = ...
   ```
   plus `close()` and, in `CoordinatorService.clearCoordinatorService()`:
   ```java
   if (jobHistoryService != null) {
       jobHistoryService.close();
   }
   ```
   
   **Runtime flow (verified directly against the current head, not assumed):**
   ```text
   Normal path (single-threaded scheduler tick, no race):
   masterActiveListener (1-thread ScheduledExecutorService, 100ms tick)
     -> CoordinatorService.checkNewActiveMaster()            
CoordinatorService.java:1133
          node loses master role
          -> clearCoordinatorService()                        
CoordinatorService.java:1150/1157, synchronized
               executorService.shutdownNow() + awaitTermination(20s)   
:1189-1201
               -> jobHistoryService.close()                    :1211-1213 (new)
                    -> removeEntryListenerQuietly x3            
JobHistoryService.java:184-195 -> :450-457
               -> resourceManager (capture-and-null) .close()  :1216-1219
               -> eventProcessor (capture-and-null) .close()   :1222-1228
   
   Node shutdown path (the one narrow exception, see Issue 1):
   CoordinatorService.shutdown()                                :1966
     masterActiveListener.shutdown()   (non-blocking on an in-flight tick)  
:1969
     clearCoordinatorService()         (called directly from the shutdown 
caller's thread)  :1977
     awaitSchedulerTermination(...)    (only awaited AFTER 
clearCoordinatorService runs)     :1978
   ```
   
   **Key Findings:**
   - I independently re-verified the core fix against the current code (not 
just the diff): all three listener ids are captured as `final UUID` fields 
(`JobHistoryService.java:122/128/134`), `close()` removes exactly those three 
via `removeEntryListenerQuietly` (`:184-195`), and the helper swallows 
exceptions per-map so a transient failure removing one listener cannot block 
the other two (`:450-457`). This matches the leak description in issue #11807 
exactly and is the correct fix for the normal, scheduler-driven transition path.
   - `clearCoordinatorService()` is `public synchronized void` 
(`CoordinatorService.java:1172`); `initCoordinatorService()`, which reassigns 
the `jobHistoryService` field (`:529`), is **not** synchronized on the same 
monitor and the field is not `volatile`. Under the routine path both are only 
ever invoked from the single `masterActiveListener` scheduler thread via 
`checkNewActiveMaster()`, so they can't interleave there. But `shutdown()` 
(`:1966-1979`) calls `masterActiveListener.shutdown()` — which does not 
interrupt or wait for a tick already in flight — and then calls 
`clearCoordinatorService()` directly from the calling thread before ever 
awaiting scheduler termination. If a `checkNewActiveMaster()` tick is 
mid-flight on the scheduler thread at that exact moment, between 
`coordinatorServiceCleared.set(false)` (`:516`, which re-arms 
`clearCoordinatorService()`'s idempotency guard) and the `jobHistoryService = 
new JobHistoryService(...)` assignment (`:529`), the shutdown-calle
 r thread's unsynchronized read of `jobHistoryService` can observe the freshly 
published instance and close **that** one instead of the outgoing generation's. 
This is a real, code-verifiable data race (two threads touching an 
unsynchronized, non-volatile field with no shared lock), not a theoretical 
concern — I traced every step above directly in the current source rather than 
trusting the earlier description of it.
   - The three-line block that deregisters `jobHistoryService` does **not** 
follow the capture-and-null idiom already used three lines below it for 
`resourceManager` (`CoordinatorService.java:1216-1219`, `ResourceManager 
manager = resourceManager; resourceManager = null;`) — that idiom exists 
precisely to make this class of hazard impossible for other fields in the same 
method, and the new code doesn't reuse it.
   - `JobHistoryService`'s constructor still performs all three 
`addEntryListener` calls unguarded (`:154-160`); if the second or third throws 
mid-construction, the instance is never published and any id(s) already 
obtained on the earlier map(s) are permanently orphaned. This is not a 
regression (pre-fix, every construction leaked all three unconditionally), but 
it's a residual gap in the same leak class the PR targets.
   - The class does not declare `implements AutoCloseable` even though it now 
has a `close()` matching that signature exactly (`:69`), and the test-only 
accessor `getEntryListenerRegistrationIds()` (`:202`) is package-private but 
not marked `@VisibleForTesting`, despite that annotation already being used 
elsewhere in this module (`CoordinatorService.java`, 
`CheckpointCoordinator.java`, `PhysicalVertex.java`, `JobMaster.java`).
   
   ## 0.5 Incorporating other reviewers
   `@SEZ9` posted a substantive review on 2026-08-21T06:54:03Z with 8 numbered 
findings. I re-verified each against the current head myself rather than 
trusting the summary:
   - Issue 1 (mutable-field race on 
`clearCoordinatorService()`/`initCoordinatorService()`): confirmed 
independently above — real, Medium. Agree with `@SEZ9`.
   - Issue 2 (closed-but-non-null field could be reused by a future refactor): 
I checked — `initCoordinatorService()` today always does `jobHistoryService = 
new JobHistoryService(...)` unconditionally, with no `if (jobHistoryService == 
null)` guard anywhere, so this is not an active bug today, only a 
forward-looking fragility. I rate it Low rather than `@SEZ9`'s Medium; `@SEZ9`, 
if you disagree I'm happy to hear the concrete path where it bites today.
   - Issue 3 (partial-construction orphan): confirmed the code path is real, 
but since it's strictly better than the pre-fix baseline (which leaked all 
three unconditionally, every time) and requires a rare exception 
mid-construction, I rate it Low rather than Medium — happy to be overruled if 
there's a concrete production trigger I'm missing.
   - Issues 4-7 (test try/finally hygiene, swallowed removal-failure 
diagnostics, `AutoCloseable`, `@VisibleForTesting`): confirmed all four 
directly against the current code (test file line 77 `liveService` has no 
try/finally; `removeEntryListenerQuietly` discards the boolean from 
`removeEntryListener`; the class has no `implements AutoCloseable`; the 
accessor has no `@VisibleForTesting`). Agree, Low severity on all four, +1 to 
`@SEZ9`.
   - Issue 8 (raw-type test class): checked the surrounding test suite's actual 
convention — the closest siblings in purpose 
(`CoordinatorServiceJobCleanupTest`, 
`CoordinatorServiceWithCancelPendingJobTest`, 
`CoordinatorServicePipelineCleanupTest`, all `CoordinatorService`-lifecycle 
tests like the new one) also extend `AbstractSeaTunnelServerTest` raw, with no 
type argument. I would not block on this — it matches its closest siblings, 
even if other, unrelated tests in the module use the generic form.
   
   None of `@SEZ9`'s findings have been addressed by a new commit — the head is 
unchanged since that review, so I'm carrying all of them forward as still-open 
below rather than treating anything as resolved.
   
   **Issue 1: `clearCoordinatorService()` can close a newer master's listeners 
if it races an in-flight `initCoordinatorService()` during `shutdown()`**
   - Location: 
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java:1211-1213`
   - Problem: see the runtime-flow trace above — the field read/write is 
unsynchronized across the two methods, and `shutdown()` gives a narrow window 
where a stale/late tick and a shutdown-triggered clear can interleave.
   - Potential risk: the *new* active master's finished-job expiration side 
effects silently stop firing cluster-wide until the next role switch, with no 
error or log line explaining why — a worse failure mode than the leak this PR 
fixes.
   - Best improvement: apply the same capture-and-clear pattern already used 
for `resourceManager` two lines below (a local variable captured under the same 
`synchronized` block, or an `AtomicReference` CAS shared with 
`initCoordinatorService()`), so `close()` can never target an instance 
published by a later activation.
   - Severity: Medium
   - Raised by another reviewer: Yes (@SEZ9)
   
   **Issue 2: closed instance left reachable via the field with no `closed` 
marker**
   - Location: `CoordinatorService.java:1211-1213`
   - Problem: nothing distinguishes a closed `JobHistoryService` from a live 
one at the type level.
   - Potential risk: purely forward-looking — only bites if a future refactor 
adds an instance-reuse branch to `initCoordinatorService()`; not reachable 
today since that method unconditionally constructs a fresh instance every time.
   - Best improvement: a `closed` boolean (cheap insurance), or fold into Issue 
5's `AutoCloseable` change.
   - Severity: Low
   - Raised by another reviewer: Yes (@SEZ9, who rated it Medium — I rate Low 
since it's not reachable in the current code path)
   
   **Issue 3: constructor can still orphan 1-2 registrations if 
`addEntryListener` throws mid-construction**
   - Location: 
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/master/JobHistoryService.java:152-160`
   - Problem: if the second/third `addEntryListener` throws (e.g. 
`HazelcastInstanceNotActiveException` racing a shutdown mid-activation), the 
constructor aborts and the instance holding the earlier id(s) is never 
published, so `close()` can never reach them.
   - Potential risk: a narrow, rare trigger; strictly better than the pre-fix 
baseline (guaranteed 3-listener leak on every transition) since at most 1-2 
registrations are orphaned on an exceptional path.
   - Best improvement: wrap the three registrations in try/catch, best-effort 
removing any already-registered ids via `removeEntryListenerQuietly` before 
rethrowing.
   - Severity: Low
   - Raised by another reviewer: Yes (@SEZ9, who rated it Medium — I rate Low 
as it's a residual gap, not a regression, and narrowly triggered)
   
   **Issue 4: `removeEntryListenerQuietly` treats the expected shutdown race 
and a genuine bug identically, and discards the boolean result of 
`removeEntryListener`**
   - Location: `JobHistoryService.java:450-457`
   - Problem: catches bare `Exception` at a single `warning` level; also 
discards the `false` return that `removeEntryListener` gives when a 
registration is already gone, so a persistent, non-exception removal failure 
would silently reintroduce the exact leak this PR targets.
   - Best improvement: catch `HazelcastInstanceNotActiveException` separately 
at fine/debug level, log other exceptions at warning with the `registrationId`, 
consider logging on an unexpected `false` return.
   - Severity: Low
   - Raised by another reviewer: Yes (@SEZ9)
   
   **Issue 5: `JobHistoryService` doesn't declare `implements AutoCloseable`**
   - Location: `JobHistoryService.java:69`
   - Problem: the class now owns a `close()` lifecycle matching 
`AutoCloseable`'s signature exactly, but doesn't declare it.
   - Best improvement: declare `implements AutoCloseable` — free, makes the 
"must be closed" contract compiler/IDE-discoverable, and lets the new test use 
try-with-resources (also fixing Issue 7 below for free).
   - Severity: Low
   - Raised by another reviewer: Yes (@SEZ9)
   
   **Issue 6: test-only accessor `getEntryListenerRegistrationIds()` isn't 
marked `@VisibleForTesting`**
   - Location: `JobHistoryService.java:202`
   - Problem: package-private, Javadoc says "only intended for tests," but 
nothing enforces that; `@VisibleForTesting` is already an established 
convention in this exact module.
   - Best improvement: add the annotation.
   - Severity: Low
   - Raised by another reviewer: Yes (@SEZ9)
   
   **Issue 7: `testCloseRemovesFinishedJobEntryListeners`'s `liveService` (and 
the create/close loop) register listeners on the shared test node without 
try/finally**
   - Location: 
`seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/master/JobHistoryServiceListenerCleanupTest.java:77-104`
   - Problem: I confirmed directly — line 77's `liveService` has no surrounding 
try/finally (the file's only try/finally block is elsewhere, guarding the 
isolated-node test). A failed assertion between creation and cleanup leaves 
live registrations on the shared node's cluster-wide IMaps for the rest of the 
JVM's test run, reproducing this PR's exact leak pattern inside the test suite 
itself and making unrelated failures in sibling tests harder to diagnose.
   - Best improvement: wrap `liveService` and each loop-created instance in 
try/finally calling `close()` (documented idempotent, so double cleanup is 
safe).
   - Severity: Low
   - Raised by another reviewer: Yes (@SEZ9)
   
   ## 1.2 Compatibility Impact
   Fully compatible. Purely additive: three new `private final UUID` fields, 
one new `public void close()`, one new package-private test accessor, and one 
new call site in an existing method. No config option, public API signature, 
checkpoint/savepoint format, or wire-protocol change. `close()` is idempotent 
and cannot throw into the caller.
   
   ## 1.3 Performance / Side-Effect Analysis
   Negligible. `close()` does at most three `IMap.removeEntryListener` calls 
once per master-role transition — transitions are rare relative to normal 
operation, not per-job or per-event. No new locks, threads, or hot-path impact.
   
   ## 1.4 Error Handling and Logging
   Adequate for the stated best-effort contract, with the diagnosability gap 
noted as Issue 4 above (bare `Exception` catch, discarded boolean result).
   
   # 2. Code Quality Assessment
   ## 2.1 Coding Standards
   New fields/methods and the test class carry Javadoc explaining purpose and 
lifecycle risk. No missing-comment gaps on the core methods. Issues 5/6 above 
are the only structural/annotation gaps.
   
   ## 2.2 Test Coverage and Test Stability
   Two new regression tests cover both the direct `close()` path and the 
production `clearCoordinatorService()` call path, using 
`Awaitility.await().ignoreExceptions().atMost(30, 
TimeUnit.SECONDS).untilAsserted(...)` for the isolated-node activation wait — 
no `Thread.sleep`, no fixed-delay guessing. The isolated second Hazelcast node 
uses a reserved port range (`TestUtils.getAvailablePort(100)`) rather than a 
default/shared port, avoiding cross-test-class port collisions in the same JVM. 
**Rating: Stable**, with the one test-hygiene caveat in Issue 7 (no try/finally 
on the shared-node registrations).
   
   ## 2.3 Documentation Updates
   Not applicable — no user-facing config, API, or behavior contract changes; 
the PR description correctly marks no documentation update as required.
   
   # 3. Architectural Soundness
   ## 3.1 Elegance of the Solution
   Precise fix for the reported leak on its primary (scheduler-driven) path. 
Minimal and additive; follows the file's own established idioms almost 
everywhere except the one place that matters most for correctness under the 
shutdown race (Issue 1).
   ## 3.2 Maintainability
   Good — `close()` and the new fields make the listener-ownership contract 
explicit where none existed before. Issues 2/5/6 are all about making that 
contract harder to violate accidentally in the future.
   ## 3.3 Extensibility
   No concerns; the pattern generalizes cleanly if a fourth cluster-wide 
listener is ever added.
   ## 3.4 Historical-Version Compatibility
   No serialized state, checkpoint format, or config option touched. A rolling 
upgrade with old and new binaries mixed is safe: old nodes keep leaking as 
before until upgraded, new nodes stop leaking, and there is no cross-version 
wire-format interaction.
   
   # 4. Issue Summary
   | # | Issue | Location | Severity |
   |---|---|---|---|
   | 1 | `clearCoordinatorService()` can close a newer master's listeners if it 
races an in-flight `initCoordinatorService()` during `shutdown()` | 
`CoordinatorService.java:1211-1213` | Medium |
   | 2 | Field left pointing at a closed instance with no `closed` marker 
(forward-looking, not an active bug today) | 
`CoordinatorService.java:1211-1213` | Low |
   | 3 | Partial-construction failure can orphan 1-2 listener registrations 
with no owner (not a regression) | `JobHistoryService.java:152-160` | Low |
   | 4 | `removeEntryListenerQuietly` treats expected-shutdown and genuine 
failures identically, discards `removeEntryListener`'s boolean result | 
`JobHistoryService.java:450-457` | Low |
   | 5 | `JobHistoryService` should declare `implements AutoCloseable` | 
`JobHistoryService.java:69` | Low |
   | 6 | Test-only accessor should be `@VisibleForTesting` | 
`JobHistoryService.java:202` | Low |
   | 7 | `testCloseRemovesFinishedJobEntryListeners` lacks try/finally cleanup 
on the shared test node | `JobHistoryServiceListenerCleanupTest.java:77-104` | 
Low |
   
   CI note: `Build` was red on this head as of yesterday's check; I traced the 
two failing jobs back then (`unit-test (8, windows-latest)` — Maven 
dependency-resolution connection reset for `seatunnel-transforms-v2`, before 
any engine-server test runs; `all-connectors-it-1 (ubuntu-latest)` — the known 
`CouchbaseIT` container-startup flake that's been hitting `dev` since 
2026-08-15) and both are infra-level, unrelated to this diff. Nothing has 
changed on the CI front since then that I can see from this pass; I'd still 
recommend rerunning the two failed jobs for a clean signal rather than treating 
the red `Build` as evidence against this PR.
   
   # 5. Merge Recommendation
   ### Conclusion: Ready to merge after fixes
   
   1. Blockers — must be fixed
      - **Issue 1** (Medium): apply the capture-and-clear idiom already used 
for `resourceManager` right below, to `jobHistoryService`, in 
`clearCoordinatorService()`, so `close()` can never target an instance 
published by a later `initCoordinatorService()` call. This is a small, 
mechanical change and it protects the exact correctness guarantee this PR 
exists to establish.
   
   2. Recommended fixes — non-blocking
      - Issues 2-7 (closed-instance marker / `AutoCloseable`, 
`@VisibleForTesting`, partial-construction cleanup, richer removal-failure 
logging, test try/finally hygiene).
   
   Overall assessment: unchanged from yesterday's round — the core fix is 
correct and well-tested for the normal, scheduler-driven master-transition path 
issue #11807 describes. `@SEZ9`'s findings remain legitimate and unaddressed 
since no new commit has landed; Issue 1 in particular is cheap enough to fix 
now that I'd rather land it in this same PR than merge with a known (if narrow) 
concurrency hazard in code whose entire purpose is closing a 
lifecycle-correctness gap. Issues 2-7 are genuine but low-cost hardening that 
can reasonably follow in this PR or a fast follow-on without blocking the 
primary fix much longer.
   


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