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

   # What Problem Does This PR Solve?
   
   `JobHistoryService` registers three cluster-wide Hazelcast 
`EntryExpiredListener`s in its constructor but never deregistered them, so 
every active-master transition leaked a full set of listeners. This PR captures 
the three registration ids and removes them from 
`CoordinatorService.clearCoordinatorService()`.
   
   Concretely what accumulates: `CoordinatorService.initCoordinatorService()` 
constructs a brand-new `JobHistoryService` every time this node becomes the 
active master. That constructor registers three entry listeners on three 
cluster-wide IMaps (finished-job state, finished-job metrics, finished-job 
vertex info). Before this PR the return value of `addEntryListener(...)` was 
discarded, and no code path ever called `removeEntryListener`. Because both 
listener classes are non-static inner classes, each abandoned registration 
keeps its enclosing `JobHistoryService` — and transitively the node engine, the 
IMaps, and the running/pending job maps — reachable from the IMap listener 
registry, and one of the listeners is the sole producer of a cleanup RPC in the 
codebase, so leaked registrations also duplicate that RPC fan-out on every real 
expiration event.
   
   Fix approach: store the three UUIDs returned by `addEntryListener` in new 
final fields, add a `close()` that deregisters exactly those three via a 
fail-soft helper, and invoke it from 
`CoordinatorService.clearCoordinatorService()`. `jobHistoryService` is 
deliberately not nulled, so read paths keep working after the role switch.
   
   # What Changed Since My Last Review
   
   Nothing in the code. I re-fetched `upstream/dev`, diffed the current head 
against the merge-base, and confirmed the production diff is byte-for-byte 
identical to what I reviewed on 2026-08-16 (the same 3 files changed). The only 
commit added since my last comment is an empty "retrigger build" commit — `git 
show --stat` on it lists no files. None of my six previously-raised issues (2 
Medium, 4 Low) have been addressed in source. I re-verified this directly 
against the current head files rather than trusting the diff being empty as a 
shortcut.
   
   # 1. Code Change Review
   
   ## 1.1 Core Logic Analysis
   
   I re-read `JobHistoryService.java` and `CoordinatorService.java` in full at 
the current head (not the diff) and re-independently re-verified the four 
load-bearing claims from my prior round, because a re-review that just restates 
prior conclusions without re-checking is worthless:
   
   - **All three listeners captured and removed** — confirmed line-for-line 
against current head.
   - **Single production construction site** — remains the only `new 
JobHistoryService(` in non-test code.
   - **CAS guard re-arms correctly** — the guard flag is reset as the first 
statement of `initCoordinatorService()`, and the CAS gate in 
`clearCoordinatorService()` is unchanged — every active/inactive cycle still 
performs exactly one cleanup.
   - **`close()` call site unchanged** — still inside the `synchronized` 
`clearCoordinatorService()`, still reached from role-loss, init-failure-retry, 
and shutdown.
   
   **Runtime path (unchanged from prior round, re-verified against current line 
numbers):**
   
   ```
   Case A  leave active master
     -> clearCoordinatorService()  [synchronized]
          CAS(false->true) guard
          awaitTermination(20s)
          jobHistoryService.close()          <-- THE FIX
                  -> removeEntryListenerQuietly x3
          resourceManager.close() ; eventProcessor.close()
   
   Case B  become active master
     -> initCoordinatorService()
          coordinatorServiceCleared.set(false)   <-- re-arms
          jobHistoryService = new JobHistoryService(...)  -> 3 fresh ids
   
   Case C  init throws  -> clearCoordinatorService(), guard is false, close() 
runs
   Case D  shutdown()   -> clearCoordinatorService() -> same close() path
   ```
   
   **Normal job lifecycle does not reach this code, and that is correct.** 
Submit -> run -> finish never touches `clearCoordinatorService()`; the removal 
fires only on master-role loss, init failure, and shutdown — not on job 
completion despite the PR title's phrasing. This is the correct scope: 
cancelling a job must not tear down the master's history listeners, and I 
confirmed no other code path calls `close()`.
   
   ## 1.2 Compatibility Impact
   
   Fully backward compatible — unchanged conclusion, re-verified. No `Option` 
added/removed/renamed. `close()` and the new id-accessor are purely additive, 
the latter package-private. `JobHistoryService` is not `Serializable`; the 
fields that are actually serialized are untouched. No checkpoint/savepoint/IMap 
payload format is affected. Rolling upgrade and downgrade are both safe — a 
downgrade simply restores the pre-fix leak, no data-format break.
   
   ## 1.3 Performance / Side-Effect Analysis
   
   This is the point of the fix, so I re-verified rather than restated:
   
   - **Net positive, confirmed:** listener count is bounded at 3 regardless of 
failover count; duplicate cleanup-RPC fan-out and duplicate cleanup-counter 
increments stop after the first master switch post-fix.
   - **New blocking work inside `synchronized clearCoordinatorService()`:** 
unchanged from prior round. Three additional synchronous 
`IMap.removeEntryListener` calls now sit inside the same synchronized method 
that already blocks on a 20s `awaitTermination` and on 
`resourceManager.close()`. Same risk class as those neighbors, not a new class 
of risk. Not a blocker.
   - **Handover window with zero registered listeners (raised in my prior 
round, still not addressed in code or Javadoc):** because 
`addEntryListener(listener, true)` is cluster-wide, there is a window between 
the outgoing master's `close()` and the incoming master's 
`initCoordinatorService()` where no listener is registered on the finished-job 
IMaps at all. An expiry landing in that window is dropped permanently — no 
reconciliation sweep exists to recover it. I re-checked whether anything since 
has added a periodic reconciliation pass — it has not. Risk remains Low 
(bounded by the up-to-20s `awaitTermination` before `close()` runs, and 
Hazelcast TTL expiry being lazy), and I still do not consider it a merge 
blocker, but it remains completely undocumented in the `close()` Javadoc, which 
a future maintainer reading only that Javadoc would not learn about.
   
   ## 1.4 Error Handling and Logging
   
   Unchanged and still correct: the removal helper catches `Exception`, logs at 
`warning` with the map name, never propagates. The null-guard on 
`jobHistoryService` covers the first-ever clear before any init. No sensitive 
data logged.
   
   No new issues found in this pass beyond the two carried over from the prior 
round:
   
   **Issue 1** *(carried over, unresolved)*
   - **Location:** `JobHistoryService` constructor; contrast the safer pattern 
used by a sibling state-store class in this module.
   - **Problem:** The three ids are `final` fields assigned inline. If the 
second or third `addEntryListener` call throws, the already-registered 
listener(s) become permanently unremovable — the exact leak this PR fixes, in 
one uncovered corner. Because `initCoordinatorService()` failure is retried 
every ~100ms by the active-master check, a persistent partial-construction 
failure would leak on every tick.
   - **Risk:** Narrow and pre-existing in shape (this exact hole existed before 
the PR too, just without any cleanup path at all), not introduced by this PR, 
but not closed by it either.
   - **Severity:** Low. Not required for merge.
   
   **Issue 2** *(carried over, unresolved)*
   - **Location:** `CoordinatorService.clearCoordinatorService()` interacting 
with `JobHistoryService`'s registration/close methods.
   - **Problem:** Zero-listener handover window described in 1.3, undocumented 
in the `close()` Javadoc.
   - **Risk:** Low probability, orphaned log files on member nodes if an expiry 
lands exactly in the window.
   - **Best improvement:** No code change required for merge; add a sentence to 
the `close()` Javadoc so the window is a documented, intentional trade-off 
rather than something the next maintainer has to rediscover from first 
principles.
   - **Severity:** Low.
   
   # 2. Code Quality Assessment
   
   ## 2.1 Coding Standards
   
   Unchanged and re-verified: ASF license header present and correct on the new 
test file. All three new fields and all three new methods carry multi-line 
Javadoc explaining purpose, not just type. No wildcard imports, no 
`System.out.println`, no single-line Javadoc.
   
   ## 2.2 Test Coverage and Test Stability — does the test actually detect a 
leak, and is it CI-safe?
   
   **Yes, it genuinely detects a leak, not just "no exception thrown."** I 
re-read the test file at the current head end-to-end. Both tests assert on 
`IMap.removeEntryListener(UUID)`'s boolean return, which is a real observation 
of the listener registry: `true` means the registration is live, `false` means 
it is gone. One test establishes a positive control first — proving the exposed 
ids are real live registrations — before running three create/close cycles and 
asserting `false` after each, which directly emulates repeated 
master-transition churn. The other test exercises the real production entry 
point (`clearCoordinatorService()`) rather than calling `close()` directly, 
which is the correct level to test at.
   
   **Both CI-stability defects I flagged in my prior round are still present, 
verbatim, in the current head — I re-read the file line-by-line to confirm, not 
just trusted my earlier finding:**
   
   - **Issue 3 (Medium, carried over, unresolved).** The test acquiring the 
coordinator service has no `await()` wrapper around it. I re-read the accessor 
directly at the current head: if the node isn't master yet, it throws 
immediately with zero retries; if it is master, the internal retry loop is 
bounded to 1.5s before throwing. The module's own sibling test file establishes 
the answer to exactly this problem — an `await().atMost(60, 
TimeUnit.SECONDS).untilAsserted(...)` idiom used repeatedly throughout that 
file. The new test does not follow this established idiom. On a loaded CI 
runner, single-node master election plus `initCoordinatorService()` completing 
inside 1.5s is not guaranteed, and this project already has a documented 
flaky-engine-server-test burden.
   - **Issue 4 (Medium, carried over, unresolved).** The isolated Hazelcast 
instance in the new test is built with only the cluster name overridden, 
inheriting default network/port settings. I re-confirmed the sibling test file 
has a private helper used at 8+ call sites specifically to avoid port-join 
contention between concurrently-running node instances in the same CI job. This 
test class also stands up a second node alongside the base-class node without 
that same protection.
   
   Both are exactly as I described them in my prior round; no attempt has been 
made to address either, and their fix cost is genuinely small (three lines 
each, using patterns that already exist verbatim elsewhere in the same test 
package).
   
   **Coverage gaps also carried forward, unaddressed:** no absolute 
listener-count assertion, no cross-instance isolation assertion (an 
over-removing `close()` that killed a concurrent instance's listeners would 
still pass this suite), no coverage of the shutdown path or init-failure path.
   
   **Stability rating: acceptable but should be hardened before merge — 
unchanged from my prior conclusion.** The fundamentals (no `Thread.sleep`, no 
floating-point tolerance, no inter-test ordering dependency, `finally`-guarded 
teardown) are genuinely good. But Issues 3 and 4 are real, 
established-pattern-diverging test-stability risks being introduced into a 
module with a known flaky-test history, and they remain unfixed in this round.
   
   ## 2.3 Documentation Updates
   
   None required, none made — correct, unchanged.
   
   # 3. Architectural Soundness
   
   ## 3.1 Elegance of the Solution
   
   Unchanged: a precise root-cause fix restoring the missing 
register/deregister symmetry at the one construction site and one teardown 
site, with no new abstraction or config. Almost entirely additive.
   
   ## 3.2 Maintainability
   
   Unchanged: `close()` sits alongside the existing resource-teardown idiom in 
`clearCoordinatorService()` and reads as the established pattern for 
coordinator-owned-resource teardown.
   
   ## 3.3 Extensibility
   
   Unchanged, non-blocking: the three `(IMap, UUID, storeName)` triples are 
three parallel fields rather than one collection. Fine at the current fixed 
count of three; worth generalizing only if a fourth listener is added later.
   
   ## 3.4 Historical-Version Compatibility
   
   Unchanged: registration ids are ephemeral in-memory state, never persisted, 
never part of a checkpoint/savepoint payload. Rolling upgrade and downgrade are 
both safe.
   
   # 4. Issue Summary
   
   | # | Issue | Location | Severity | Status |
   |---|-------|----------|----------|--------|
   | 1 | Constructor assigns three listener ids to `final` fields inline; a 
partial-construction failure leaks the already-registered listener(s) 
unremovably | `JobHistoryService.java` (constructor) | Low | Carried over, 
unresolved |
   | 2 | Zero-listener handover window is undocumented in the `close()` Javadoc 
| `CoordinatorService.java`, `JobHistoryService.java` | Low | Carried over, 
unresolved |
   | 3 | Coordinator-service accessor in the new test has no Awaitility guard, 
diverging from the established sibling-test pattern used at 8+ sites in the 
same module | `JobHistoryServiceListenerCleanupTest.java` | Medium | Carried 
over, unresolved |
   | 4 | Isolated Hazelcast instance in the new test built from default network 
config, unlike the port-contention-safe pattern used throughout the sibling 
test file | `JobHistoryServiceListenerCleanupTest.java` | Medium | Carried 
over, unresolved |
   | 5 | No cross-instance isolation assertion and no absolute listener-count 
assertion | `JobHistoryServiceListenerCleanupTest.java` | Low | Carried over, 
unresolved |
   | 6 | Non-synchronized `initCoordinatorService()` can theoretically register 
fresh listeners after the shutdown thread already ran `close()`; pre-existing 
structural asymmetry, not introduced by this PR | `CoordinatorService.java` | 
Low | Carried over, out of scope for this PR |
   
   No new issues found in this re-review beyond the six carried over from my 
last round.
   
   # 5. CI Status
   
   Re-checked live at review time, not inferred from a stale snapshot: the 
apache-side pointer showed `FAILURE`, but per the standing project fact that 
the apache-side check is only a pointer, the real signal is the fork run. I 
pulled the full job list and the failing job's raw log: `unit-test` passed on 
all four matrix legs (the module most directly exercising this change), 
`engine-v2-it` passed on both JDKs. The single failure is one connector-IT job, 
and I pulled the actual Maven output for it: the failure is a Testcontainers 
container-startup failure for a connector this PR does not touch (the diff is 
confined to `seatunnel-engine`). This reads as an environment/infra flake, not 
a regression caused by this change, but per project policy the PR still needs a 
green CI run before merge; this failure has not yet been re-verified as 
transient by a rerun.
   
   # 6. Merge Recommendation
   
   ### Conclusion: Not ready to merge. Source logic remains correct; CI is red 
on an apparently-unrelated infra flake, and two Medium test-stability issues 
from my prior review are still unaddressed one round later.
   
   1. **Blockers — must be fixed before merge**
      - **CI is red.** The fork run's connector-IT job failed on a 
Testcontainers startup issue, unrelated to this diff on its face, but per 
project policy an unverified/failing CI blocks merge regardless of cause. This 
needs a rerun to confirm it is transient before this PR can be considered 
CI-clean.
      - **Issue 3 (Medium):** wrap the coordinator-service acquisition in the 
new test in the same `await().atMost(...).untilAsserted(...)` idiom already 
used 8+ times in the sibling test file. This is the most likely source of a 
future flaky-CI report on this exact new test, and remains a three-line fix not 
yet made across two review rounds.
      - **Issue 4 (Medium):** build the isolated instance in the new test from 
an explicit config with an allocated port and raised join-port-try-count, 
mirroring the sibling test's existing helper.
   
   2. **Recommended fixes — non-blocking**
      - Issue 1 (Low): make the three constructor registrations all-or-nothing, 
or adopt the volatile-field + null-guard idiom the sibling state stores already 
use.
      - Issue 2 (Low): document the zero-listener handover window in the 
`close()` Javadoc.
      - Issue 5 (Low): add a cross-instance isolation assertion.
      - Issue 6 (Low): out of scope, note only.
   
   **Overall assessment.** The production fix itself remains correct, minimal, 
and well-reasoned on this third independent pass: I re-derived the leak 
mechanism, re-confirmed the CAS re-arm, re-confirmed there is exactly one 
production construction site and one teardown site, and re-confirmed 
compatibility is clean on every axis (no API, no config, no 
serialized/checkpoint state). Nothing about the source correctness conclusion 
has changed since my last round, and I found no new production-code issues in 
this pass. What has not changed is more important right now: this PR sat for a 
full day with CI queued, then went green everywhere that matters to the actual 
bug except a plausibly-unrelated Testcontainers flake, and neither of the two 
Medium test-stability findings from my last review — both three-line fixes, 
both following patterns that already exist in the very same test file's sibling 
class — were acted on before the retrigger. Holding my own PR to the full bar 
means not lett
 ing "CI will probably go green on rerun" substitute for actually fixing 
findings I already wrote down and already know how to fix. I'd fix Issues 3 and 
4 now, rerun the failed job, and only then move this to mergeable.
   
   (Note: since I authored this PR myself, this review is posted as a plain 
comment rather than a formal approve/request-changes review — GitHub doesn't 
allow self-approval, and I want the same level of independent scrutiny applied 
here as to anyone else's PR.)
   


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