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

   What changed since my last review: nothing in the code. I re-fetched the 
current head (`fbe4e66d2df`) and diffed it against `dev`'s merge-base — the 
production diff is byte-for-byte identical to what I reviewed on 2026-08-17 
(same 3 files, +253/-6). The only new commit is another empty `[Chore][Zeta] 
Retrigger PR #11809 build` commit (`git show --stat` on it lists no files). I 
re-read `JobHistoryService.java`, `CoordinatorService.java`, and the new test 
file end-to-end at the current head rather than trusting that "diff unchanged" 
claim as a shortcut, since this round specifically asks me to re-confirm the 
"3, not 1" listener-completeness claim from issue #11807 and check for a 
removal-before-drain race.
   
   # Re-verification against issue #11807's "3 listeners, not 1" claim
   
   Confirmed directly at the current head, not inferred from the PR description:
   
   - `JobHistoryService` registers exactly three cluster-wide 
`EntryExpiredListener`s in its constructor, one each on `finishedJobStateImap`, 
`finishedJobMetricsImap`, `finishedJobDAGInfoImap` 
(`JobHistoryService.java:154-162`), and captures all three returned `UUID`s 
into `final` fields (`:117-134`).
   - `close()` (`:184-195`) calls `removeEntryListenerQuietly()` three times, 
once per captured id, covering all three maps — not a subset.
   - `removeEntryListenerQuietly()` (`:450-457`) wraps each 
`imap.removeEntryListener(id)` in try/catch, logging at `warning` on failure, 
never throwing — so a transient failure removing one listener cannot prevent 
the other two from being attempted, and cannot propagate out of `close()`.
   - A repo-wide sweep for `addEntryListener`/`addLocalEntryListener` confirms 
these are the only three registrations tied to `JobHistoryService`'s lifecycle; 
the two other engine-side `addEntryListener` call sites (in the 
checkpoint/metrics state stores) already have their own correct capture/remove 
symmetry and are unrelated to this fix.
   
   So: all three leaking listeners identified in #11807 are captured and 
removed by this PR, not just one. No partial-removal or double-removal issue — 
`IMap.removeEntryListener` on an already-removed id returns `false` rather than 
throwing, so calling `close()` twice (which can't currently happen, since it's 
gated by the CAS below, but is worth noting as safe-by-construction) would be a 
no-op the second time.
   
   # Re-verification against the "removal-before-drain race" / 
abnormal-termination question
   
   - The single production teardown call site is 
`CoordinatorService.clearCoordinatorService()` 
(`CoordinatorService.java:1211-1213`, `if (jobHistoryService != null) { 
jobHistoryService.close(); }`), which is itself `synchronized` and gated by 
`coordinatorServiceCleared.compareAndSet(false, true)` (`:1174`). 
`initCoordinatorService()` resets that flag as its first statement (`:516`), so 
every active→inactive cycle re-arms the guard and `close()` runs exactly once 
per cycle — not zero times, not twice.
   - `clearCoordinatorService()` is reached on **all three** relevant paths, 
not only normal completion: role loss (`checkNewActiveMaster` Case A), 
`initCoordinatorService()` throwing during startup (Case C, caught at 
`:1152-1160` — the guard being freshly reset at `:516` means a failed partial 
init still gets cleaned up), and full node `shutdown()` (`:1966-1977`). So 
listener deregistration fires on abnormal termination and failure paths too, 
not just the happy path.
   - Hazelcast does not quiesce in-flight `entryExpired()` callbacks on 
`removeEntryListener()`, so an event that started dispatching just before 
`close()` can still complete concurrently. That's benign here: 
`incrementFinishedJobCleanupTotal` uses a `ConcurrentHashMap<String, 
AtomicLong>` (`:115` / `:440-442`), and the `CleanLogOperation` fan-out is 
already wrapped in its own try/catch, so a concurrent in-flight callback can't 
corrupt state or throw out of the listener thread.
   - The one genuine timing gap — which I flagged in my 08-16 round and is 
still accurate and still undocumented — is the opposite direction: a window 
with **zero** registered listeners between the outgoing master's `close()` and 
the incoming master's fresh registrations in `initCoordinatorService()`. An 
expiry landing in that exact window is silently dropped (no reconciliation 
sweep exists). This is Low severity (bounded by the up-to-20s 
`awaitTermination` that runs before `close()`, and Hazelcast TTL expiry being 
lazy), not a correctness blocker, but the `close()` Javadoc still doesn't 
mention it.
   
   # Current issue list (unchanged from my 2026-08-17 review — carried forward, 
not re-numbered)
   
   | # | Issue | Location | Severity | Status |
   |---|-------|----------|----------|--------|
   | 1 | Constructor assigns the three listener ids to `final` fields inline; a 
partial-construction failure (2nd/3rd `addEntryListener` throwing) leaks the 
already-registered listener(s) unremovably | `JobHistoryService.java` 
constructor (`:154-162`) | Low | Unresolved, 3rd round |
   | 2 | Zero-listener handover window (described above) is undocumented in the 
`close()` Javadoc | `CoordinatorService.java:1211-1213`, 
`JobHistoryService.java:184-195` | Low | Unresolved, 3rd round |
   | 3 | `testClearCoordinatorServiceDeregistersJobHistoryListeners` calls 
`coordinatorService.getJobHistoryService()`/`seaTunnelServer.getCoordinatorService()`
 with no `await()` wrapper. I re-read `SeaTunnelServer.getCoordinatorService()` 
at the current head (`SeaTunnelServer.java:286-317`): if the node isn't master 
yet it throws immediately with zero retries; if it is master but the 
coordinator hasn't finished initializing, the internal retry is capped at 
3×500ms=1.5s before throwing. The sibling test file already establishes the 
correct idiom for this exact call (`await().atMost(60, 
TimeUnit.SECONDS).untilAsserted(...)`, used 8+ times in 
`CoordinatorServiceTest.java`), and the new test doesn't follow it | 
`JobHistoryServiceListenerCleanupTest.java:104-119` | Medium | Unresolved, 3rd 
round |
   | 4 | The isolated Hazelcast instance in the new test is created via 
`SeaTunnelServerStarter.createHazelcastInstance(clusterName)` with only the 
cluster name overridden — default network/port config. The sibling test file 
has an existing helper, 
`createHazelcastInstanceWithJoinPortTryCount(clusterName, 100)`, specifically 
to avoid join/port contention when multiple node instances run in the same CI 
process; this test (which extends `AbstractSeaTunnelServerTest` and therefore 
already has a base-class node running) doesn't use it | 
`JobHistoryServiceListenerCleanupTest.java:105-108` | Medium | Unresolved, 3rd 
round |
   | 5 | Neither test asserts a cross-instance isolation case (closing one 
`JobHistoryService` must not remove a different, still-live instance's 
registrations) or an absolute listener-count on the maps — an over-removing 
`close()` would still pass this suite | 
`JobHistoryServiceListenerCleanupTest.java:60-94` | Low | Unresolved, 3rd round 
|
   | 6 | Non-`synchronized` `initCoordinatorService()` can theoretically 
register fresh listeners after the shutdown thread has already run `close()`. 
Pre-existing structural asymmetry, not introduced by this PR, out of scope here 
| `CoordinatorService.java:515` vs `:1172` | Low | Unresolved, out of scope |
   
   No new issues found in this pass. The production logic conclusion is 
unchanged and, on this fourth independent read, still holds: the fix is 
complete with respect to all three listeners named in #11807, is reached on the 
normal-role-loss, init-failure, and shutdown paths (not only the "happy path"), 
and does not introduce a double-removal or a removal-before-drain correctness 
bug.
   
   # CI status (checked live at review time)
   
   The apache-side `Build` check still shows `pending` — per the standing 
project fact, that's only a pointer, so I went to the actual fork run 
(`DanielLeens/seatunnel`, run for head `fbe4e66d2df`). All 30 jobs have 
completed except the run's own top-level status, which is still reporting 
`queued` — that looks like an API-side lag rather than real pending work, since 
I could not find any job still in `queued`/`in_progress` state.
   
   Results:
   - `unit-test (8, ubuntu-latest)`, `unit-test (8, windows-latest)`: pass.
   - `unit-test (11, ubuntu-latest)`: **cancelled** (not evaluable on its own; 
needs a rerun).
   - `unit-test (11, windows-latest)`: **failure** — 
`TaskExecutionServiceTest.before` → `IllegalStateException: Node failed to 
start!` after 311.37s.
   - All `all-connectors-it-*` and `jdbc-connectors-it-part-*` legs that ran: 
pass.
   
   I pulled the raw log for the failing job and specifically checked whether 
this PR's new test is implicated, since Issue 4 above is exactly about 
port/network isolation. It isn't: `JobHistoryServiceListenerCleanupTest` ran 
earlier in the same job and passed cleanly (2/2, 2.972s), its Hazelcast 
instances logged a clean `SHUTDOWN` before the next class started, and the 
failing class (`TaskExecutionServiceTest`) is unrelated code this PR doesn't 
touch, in a different package. This matches a `windows-latest` engine-server 
"Node failed to start!" flake pattern that predates this PR and has recurred on 
unrelated PRs hitting a different test class each time — it reads as 
environmental, not caused by this diff.
   
   Given project policy, an unverified/failing CI still blocks merge regardless 
of suspected cause. Recommended next step: rerun just the two JDK-11 
`unit-test` jobs (job-level rerun, not a full workflow rerun) rather than 
treating this as a code problem.
   
   # Merge recommendation
   
   ### Conclusion: Not ready to merge yet — source logic is correct and 
complete, but CI needs a clean rerun and the two Medium test-stability issues 
should be fixed now rather than carried to a 4th round.
   
   1. **Blockers — must be fixed before merge**
      - CI: `unit-test (11, windows-latest)` failed (environmental flake, 
unrelated to this diff on the evidence above) and `unit-test (11, 
ubuntu-latest)` was cancelled. Both need a clean pass before merge per project 
policy.
      - Issue 3 (Medium): add the `await().atMost(...).untilAsserted(...)` 
guard around the coordinator-service acquisition in 
`testClearCoordinatorServiceDeregistersJobHistoryListeners`, matching the 
existing idiom in `CoordinatorServiceTest.java`. Three review rounds in and 
this is still a 3-line fix using a pattern that already exists in the same test 
package.
      - Issue 4 (Medium): build the isolated Hazelcast instance from 
`createHazelcastInstanceWithJoinPortTryCount(...)` instead of the 
default-network `createHazelcastInstance(clusterName)`, matching the same 
sibling file.
   
   2. **Recommended fixes — non-blocking**
      - Issue 1 (Low): make the three constructor registrations all-or-nothing, 
or switch to the `volatile UUID` + null-guard idiom the sibling state-store 
classes already use.
      - Issue 2 (Low): document the zero-listener handover window in the 
`close()` Javadoc.
      - Issue 5 (Low): add a cross-instance isolation assertion so an 
over-removing `close()` would be caught.
      - Issue 6 (Low): out of scope for this PR, noted only.
   
   **Overall assessment.** On a fourth independent pass through the production 
code, the fix remains correct, minimal, and complete: all three listeners named 
in issue #11807 are captured and removed, the CAS guard re-arms correctly so 
cleanup runs exactly once per transition, and the single teardown call site is 
reached on role-loss, init-failure, and shutdown, not only normal completion. I 
found no removal-before-drain race and no double-removal risk — 
`IMap.removeEntryListener` is idempotent and the try/catch in 
`removeEntryListenerQuietly` prevents one failed removal from blocking the 
other two. What's holding this back is unchanged from my last round: CI isn't 
clean yet (one apparently-environmental Windows flake plus one cancelled job 
needing a rerun), and the two Medium test-stability findings — both cheap, both 
following patterns that already exist verbatim in the same test file's sibling 
class — still haven't been picked up across three rounds of review. I'd rather 
fix
  Issues 3 and 4 now and get a clean CI run than keep re-confirming the same 
two open items on a fourth pass.
   
   (Note: since I authored this PR myself, this 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