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

   Since this is my own PR, GitHub won't let me submit a formal review state on 
it, so this is posted as a plain comment again. This is round 7 on this PR (my 
sixth pass plus @SEZ9's independent review). The code itself is unchanged since 
my last comment (head is still `f953140ec`), so I'm not re-deriving the whole 
design from scratch again — I independently re-verified the current source 
against every load-bearing claim from round 6 and @SEZ9's review (all confirmed 
accurate), and I re-checked CI live rather than reusing round 6's "still 
queued" status. That live check changes the conclusion: **CI has since 
completed and the current head does not compile.** This is a hard blocker that 
supersedes my own round-6 "Ready to merge after fixes."
   
   # What Problem Does This PR Solve?
   
   When a worker leaves the Hazelcast cluster during a routine, 
operator-initiated scale-down, `CoordinatorService` fails its deployed task 
groups and `PhysicalVertex.stateProcess()` always logged that at `ERROR`, 
making expected scale-down noise indistinguishable from a real engine fault. 
After six rounds of iteration (fixing a message-construction bug that made the 
original classification a no-op, narrowing "graceful" from "any member removal" 
to a self-reported per-address marker, gating the marker on Hazelcast's 
`terminate` flag, and — following @SEZ9's finding — relocating the marker 
publication to the correct pre-shutdown SPI hook), the mechanism now publishes 
a best-effort "I'm leaving on purpose" marker into a new 
`engine_gracefulMemberRemoval` IMap before a graceful shutdown, consumes it 
when classifying the resulting task failure, and downgrades only the exact 
node-offline message to `WARN` when the marker proves the departure was 
graceful. Any unproven/abrupt departure
  (crash, OOM, forced termination, partition) still logs `ERROR`, unchanged 
from today.
   
   # 1. Code Change Review
   
   ## 1.1 Core Logic Analysis
   
   I independently re-verified the current diff against the tree directly (not 
by re-reading the prior rounds' write-ups) and confirm every load-bearing claim 
from round 6 and @SEZ9's review holds:
   
   - `SeaTunnelServer.java`: `SeaTunnelServer` now additionally implements 
`GracefulShutdownAwareService`; `onShutdown(long, TimeUnit)` calls 
`markLocalGracefulMemberRemoval()`, and the marker-write call is removed 
entirely from `ManagedService.shutdown(boolean terminate)` (confirmed by 
reading the full `shutdown()` body — every remaining line there operates on 
separately-mocked, null-guarded fields, none of them `nodeEngine`). `init()` 
clears this node's own stale marker on startup.
   - `CoordinatorService.java`: adds `gracefulMemberRemovalIMap`, wired in 
`initCoordinatorService()`; 
`buildMemberRemovedOfflineMessage`/`buildMemberRemovedFailureState`/`isGracefulMemberRemovalMarkerValid`/`consumeGracefulMemberRemovalMarker`
 are exactly as described — the graceful branch builds `TaskExecutionState` via 
the `String` constructor (plain message, regex-matchable), the non-graceful 
branch still wraps in `new JobException(...)` (full stack trace, never matches, 
`ERROR` unchanged). `failedTaskOnMemberRemoved` calls 
`consumeGracefulMemberRemovalMarker` once, up front, and threads the resulting 
boolean through both `makeTasksFailed` call sites.
   - `PhysicalVertex.java`: `DEPLOYED_NODE_OFFLINE_ERROR_PATTERN` and 
`isDeployedNodeOfflineFailure()` are exactly as described; `stateProcess()`'s 
`FAILED` case now branches `log.warn` vs `log.error` on that predicate. 
Confirmed the message-template is genuinely defined in two disconnected places 
with no shared constant (`CoordinatorService.java`'s `%s`/`%s` format string 
vs. `PhysicalVertex.java`'s hand-written regex) — @SEZ9's Issue 8 is real and I 
could find no test that feeds one into the other.
   - `Constant.java`: adds `IMAP_GRACEFUL_MEMBER_REMOVAL` and a 5-minute 
`GRACEFUL_MEMBER_REMOVAL_MARK_TTL_MILLIS`, both with adequate Javadoc this 
round (closing part of round-6's Issue 3 from earlier rounds).
   
   **Runtime path** (unchanged from round 6's trace, re-confirmed against 
current source):
   
   ```text
   Graceful path:
     operator-initiated shutdown -> Hazelcast Node.shutdown(false)
       -> callGracefulShutdownAwareServices() -> SeaTunnelServer.onShutdown()   
[SeaTunnelServer.java:233-237]
            markLocalGracefulMemberRemoval() -> 
gracefulMemberRemovalIMap.put(thisAddress, now)
     -> (network) coordinator observes memberRemoved
       -> CoordinatorService.failedTaskOnMemberRemoved(event)                  
[CoordinatorService.java:2039]
            gracefulMemberRemoval = 
consumeGracefulMemberRemovalMarker(lostAddress)  [:2008-2013]
       -> makeTasksFailed(..., gracefulMemberRemoval)
            -> buildMemberRemovedFailureState(...): plain-String 
TaskExecutionState when graceful
     -> PhysicalVertex.stateProcess() case FAILED: 
isDeployedNodeOfflineFailure(errorMsg) matches -> log.warn
   
   Abrupt path (kill -9 / OOM / partition / forced terminate):
     no marker ever written for that address -> 
consumeGracefulMemberRemovalMarker returns false
     -> buildMemberRemovedFailureState(..., false) -> 
JobException(offlineMessage) -> stack trace
     -> isDeployedNodeOfflineFailure() never matches a multi-line stack trace 
-> log.error (unchanged from dev)
   ```
   
   **What I found that neither round 6 nor @SEZ9 caught: the current head does 
not compile.** I pulled the live CI run for the current head SHA 
(`f953140ecf6ee1e96544679c0ed28d91e51d94b0`, `DanielLeens/seatunnel` fork run 
`32484015488`, completed 2026-08-22, `conclusion: failure`) rather than 
trusting round 6's "still queued at review time" note. All four `unit-test` 
matrix legs (JDK 8/11 x ubuntu/windows) fail, and I pulled the raw job log to 
find the actual cause rather than assuming it was another Windows/CDC-style 
environmental flake:
   
   ```
   [ERROR] COMPILATION ERROR :
   [ERROR] 
.../seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/SeaTunnelServerShutdownTest.java:[55,16]
 error: no suitable method found for thenReturn(IMap<Address,Long>)
   [ERROR]     method OngoingStubbing.thenReturn(IMap<Object,Object>) is not 
applicable
   [ERROR]       (argument mismatch; IMap<Address,Long> cannot be converted to 
IMap<Object,Object>)
   ```
   
   This is a `mvn test-compile` failure in `seatunnel-engine-server`, at:
   ```java
   IMap<Address, Long> gracefulMemberRemovalIMap = mock(IMap.class);   // 
raw-typed mock, unchecked cast
   ...
   when(hazelcastInstance.getMap(Constant.IMAP_GRACEFUL_MEMBER_REMOVAL))
           .thenReturn(gracefulMemberRemovalIMap);                     // javac 
can't reconcile the generics here
   ```
   `mock(IMap.class)` produces a raw `IMap`; assigning it to the parameterized 
local variable is an unchecked conversion, and when it's then passed into 
`.thenReturn(...)` on the `OngoingStubbing<IMap<Object,Object>>` inferred from 
`HazelcastInstance.getMap(String)`'s unbound type parameters, javac 8/11 both 
fail to unify `IMap<Address,Long>` with `IMap<Object,Object>`. I confirmed this 
same single root cause cascades into the entire CI matrix, not just 
`unit-test`: I independently pulled the log for `all-connectors-it-1 (8, 
ubuntu-latest)` — a job with zero topical relation to this PR — and it shows 
the identical `COMPILATION ERROR` at the identical line, because that job's 
Maven invocation also builds `seatunnel-engine-server` with `-am` in the shared 
reactor. That is why essentially the entire 60+-job CI matrix is red on this 
head, not because of 60 unrelated flakes.
   
   **Key findings:**
   - The design and runtime-path logic (the actual subject of six rounds of 
review) remain sound and unchanged — the compile break is confined to one test 
file and does not reflect a defect in the production code path.
   - This is a real, reproducible compile failure on this exact head, not a 
flake — it is 100% deterministic (a `javac` generics-inference limitation, not 
a race or environment issue), and it fully explains the "mass CI failure" 
pattern rather than requiring 60 separate unrelated-flake diagnoses.
   - The fix is small and well-understood: either give `getMap` an explicit 
type witness (`hazelcastInstance.<Address, Long>getMap(...)`), or stub via 
`doReturn(gracefulMemberRemovalIMap).when(hazelcastInstance).getMap(...)` 
(which sidesteps `OngoingStubbing`'s generic-return-type inference entirely and 
is the more common idiom elsewhere in this codebase for exactly this class of 
Mockito+generics friction).
   
   ## 1.2 Compatibility Impact
   
   Fully compatible. No config `Option`, public API, SPI contract break 
(implementing an additional Hazelcast-defined interface is additive), or 
checkpoint/wire format touched. Rolling-upgrade safety is preserved: an 
old-version member never implements `GracefulShutdownAwareService`, so its 
departure is never marked graceful and falls back to the pre-existing 
always-`ERROR` behavior identical to current `dev`.
   
   ## 1.3 Performance / Side-Effect Analysis
   
   Negligible — one `IMap.put()` on graceful shutdown (now dispatched from 
Hazelcast's own bounded `graceful-shutdown` executor rather than inline in 
`ManagedService.shutdown()`), one `IMap.remove()` per `memberRemoved` event on 
the already-slow-path of task failure handling. No new locks, no hot-path 
allocation.
   
   ## 1.4 Error Handling and Logging
   
   **Issue 1 — Current head does not compile; CI is red on a genuine compile 
error, not a flake**
   - Location: 
`seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/SeaTunnelServerShutdownTest.java:55`
   - Problem: `mock(IMap.class)` (raw type) assigned to `IMap<Address, Long>` 
and then passed to `.thenReturn(...)` fails javac's generic-return-type 
inference against `HazelcastInstance.getMap(String)`'s unbound `<K, V>` 
parameters, on both JDK 8 and JDK 11.
   - Potential risk: this is not cosmetic — it blocks every CI job in the 
matrix (confirmed identical failure in an unrelated `all-connectors-it-1` job 
via the shared `-am` reactor build), and it must be treated as a 
release-blocking correctness gate regardless of how sound the surrounding 
design is.
   - Best improvement: 
`Mockito.doReturn(gracefulMemberRemovalIMap).when(hazelcastInstance).getMap(Constant.IMAP_GRACEFUL_MEMBER_REMOVAL);`
 (avoids `OngoingStubbing`'s generic return-type inference entirely), or add an 
explicit type witness: `hazelcastInstance.<Address, Long>getMap(...)`.
   - Severity: Critical
   - Raised by another reviewer: No — CI was still queued at the time of both 
my round-6 review and @SEZ9's review; neither of us could have caught this from 
source reading alone since the code compiles under simple visual inspection 
(the type mismatch is a javac inference artifact, not a type error a human eye 
reliably catches). This is a genuine miss on my part in round 6 for asserting 
"no blockers remaining... contingent on CI" without following up once CI 
actually completed. I'm not going to let that stand uncorrected.
   
   **Issue 2 (carried over from round 6, unaddressed) — `onShutdown()`'s 
comment misstates the mechanism it relies on**
   - Location: `SeaTunnelServer.java:233-235`
   - The comment says "before the node becomes PASSIVE"; per Hazelcast 5.1's 
`Node#shutdown(boolean)`, `NodeState` is set to `PASSIVE` unconditionally 
before `callGracefulShutdownAwareServices()` runs, so the node is already 
`PASSIVE` at that point — the actual reason the placement is safe is that it 
runs before `NodeEngineImpl.shutdown()` tears down 
`proxyService`/`operationService` (which back `getMap()`). I verified this 
against the actual Hazelcast 5.1 source in round 6 and it's still uncorrected 
in the current diff.
   - Risk: low on its own; a future maintainer reasoning from the wrong premise 
could relocate this call to somewhere that again races proxy-service teardown.
   - Best improvement: reword to describe the real mechanism (service-shutdown 
ordering, not `NodeState`).
   - Severity: Low
   - Raised by another reviewer: Partially (@SEZ9 raised the underlying 
placement concern; the comment-accuracy point is mine).
   
   **Issue 3 (carried over, unaddressed) — marker consumption is destructive 
with no Hazelcast-native TTL**
   - Location: `CoordinatorService.java` (`consumeGracefulMemberRemovalMarker`, 
`remove()` before any task group is failed), `SeaTunnelServer.java` (`put()` 
with no TTL overload).
   - A master failover between `remove()` and the end of the propagation loop, 
or any departure whose marker is never read back, leaves an unbounded, 
un-evicted entry.
   - Best improvement: use the TTL-bearing `put(key, value, ttl, unit)` 
overload; consider non-destructive read until the propagation loop completes.
   - Severity: Medium
   - Raised by another reviewer: Yes (@SEZ9, Issues 2 and 4).
   
   **Issue 4 (carried over, unaddressed) — graceful branch swaps the 
`TaskExecutionState` payload from `JobException` to a bare `String`**
   - Location: `CoordinatorService.buildMemberRemovedFailureState`.
   - Downstream consumers of `throwableMsg` (job history, REST error fields) 
see a differently-shaped payload depending on which branch fired.
   - Best improvement: keep `JobException` on both branches; carry the 
classification as an explicit flag instead of a payload-shape difference.
   - Severity: Medium
   - Raised by another reviewer: Yes (@SEZ9, Issues 3 and 6).
   
   **Issue 5 (carried over, unaddressed) — offline-message template duplicated 
with no shared constant or round-trip test**
   - Location: `CoordinatorService.buildMemberRemovedOfflineMessage` (format 
string) vs. `PhysicalVertex.DEPLOYED_NODE_OFFLINE_ERROR_PATTERN` (regex).
   - I re-confirmed directly against the current tests: `PhysicalVertexTest` 
uses a hand-written literal, `CoordinatorServiceMemberRemovedTest` never feeds 
`buildMemberRemovedOfflineMessage`'s actual output into 
`isDeployedNodeOfflineFailure`. No test today would catch drift between the two.
   - Best improvement: derive both from one shared constant, or add a 
round-trip test.
   - Severity: Medium
   - Raised by another reviewer: Yes (@SEZ9, Issue 8). This is the cheapest, 
highest-leverage of the non-blocking follow-ups.
   
   **Issue 6 (carried over, unaddressed) — classification by free-text regex 
crosses a trust boundary**
   - Location: `PhysicalVertex.isDeployedNodeOfflineFailure`, fed by 
`errorByPhysicalVertex`, which can also be populated from 
worker-reported/connector-originated messages.
   - A connector exception whose message happens to match the exact template 
would be silently downgraded to `WARN`.
   - Best improvement: carry the classification as a typed flag on 
`TaskExecutionState`, not message text.
   - Severity: Medium
   - Raised by another reviewer: Yes (@SEZ9, Issue 5).
   
   **Issue 7 (carried over, unaddressed) — no docs/en / docs/zh update for the 
new IMap / WARN-vs-ERROR classification**
   - Severity: Low
   - Raised by another reviewer: Yes (@SEZ9, Issue 7).
   
   # 2. Code Quality Assessment
   
   ## 2.1 Coding Standards
   
   ASF headers present, Javadoc present on all new non-trivial fields/methods 
(a real improvement over round 2's gap), no wildcard imports, no 
`System.out.println`. The one standards issue is Issue 1 above — this specific 
test does not compile.
   
   ## 2.2 Test Coverage and Test Stability
   
   The classification logic (`isDeployedNodeOfflineFailure`), the 
message-construction branching (`buildMemberRemovedFailureState`), and the 
marker-write/no-write branching in `SeaTunnelServerShutdownTest` are all 
exercised as plain, synchronous, mock-based unit tests — no `Thread.sleep`, no 
shared static state, no timing dependence, no containers/network I/O. Absent 
the compile error, this would rate **Stable**. Because of Issue 1, **the actual 
current rating must be: the test suite does not build, and therefore provides 
zero coverage on this head** — a compile failure is a stronger negative signal 
than any flaky-test pattern, since it means none of these otherwise 
well-designed tests actually ran. I'm not going to soften this by rating the 
*design* of the tests separately from the fact that they don't compile; the 
mandatory rating for this PR as it stands is **High risk**, purely on the 
compile break, with a clear and narrow one-line fix.
   
   The remaining, pre-existing coverage gap (unchanged from round 6): no 
integration-level test exercises the real distributed `IMap` write/read round 
trip through an actual multi-node Hazelcast cluster — both my round-2 review 
and @SEZ9 flagged this; still open, non-blocking.
   
   ## 2.3 Documentation Updates
   
   Not applicable for this specific commit's production code; see Issue 7 for 
the pre-existing, PR-wide documentation gap (operator-visible WARN/ERROR 
classification change, new IMap).
   
   # 3. Architectural Soundness
   
   ## 3.1 Elegance of the Solution (Precise fix / Temporary workaround / 
Long-term solution)
   
   The design itself is a **precise fix** for the stated problem, verified 
independently across six rounds including against actual upstream Hazelcast 5.1 
source for the SPI-ordering question @SEZ9 raised. It remains fundamentally a 
message-text classification at the `PhysicalVertex` boundary (Issues 4-6), 
which is a reasonable minimal-diff tradeoff for this PR's scope, not a design 
defect. The one thing keeping this out of "long-term solution" territory today 
is purely the build break in Issue 1, which is mechanical, not conceptual.
   
   ## 3.2 Maintainability
   
   Good and improving round over round, aside from the still-open 
comment-accuracy nit (Issue 2) and the message-duplication risk (Issue 5) that 
a future edit could silently regress.
   
   ## 3.3 Extensibility
   
   Unchanged from round 6 — the self-reported marker pattern generalizes to 
future "was this expected?" classifications if needed.
   
   ## 3.4 Historical-Version Compatibility
   
   No impact — no serialized state, checkpoint, or wire format touched; old, 
marker-unaware members simply never get a marker, unchanged always-`ERROR` 
fallback.
   
   # 4. Issue Summary
   
   | # | Issue | Location | Severity | Raised by another reviewer |
   |---|---|---|---|---|
   | 1 | Current head fails `test-compile` in `seatunnel-engine-server` 
(Mockito/javac generics inference), cascading a red CI matrix across ~60 
unrelated jobs via the shared `-am` reactor build | 
SeaTunnelServerShutdownTest.java:55 | Critical | No |
   | 2 | `onShutdown()`'s comment misstates the real mechanism (claims "before 
PASSIVE"; actually before proxy/operation-service teardown, while already 
PASSIVE) | SeaTunnelServer.java:233-235 | Low | Partially (@SEZ9) |
   | 3 | Marker consumption destructive, no Hazelcast-native TTL; 
failover/never-consumed entries can leak or misclassify | 
CoordinatorService.java, SeaTunnelServer.java | Medium | Yes (@SEZ9) |
   | 4 | Graceful branch swaps `JobException` payload for a bare `String`, 
inconsistent with the sibling branch | CoordinatorService.java | Medium | Yes 
(@SEZ9) |
   | 5 | Offline-message template duplicated (format string vs. regex), no 
shared constant or round-trip test | CoordinatorService.java, 
PhysicalVertex.java | Medium | Yes (@SEZ9) |
   | 6 | Classification by free-text regex match crosses a trust boundary | 
PhysicalVertex.java | Medium | Yes (@SEZ9) |
   | 7 | No docs/en / docs/zh update for the new IMap / WARN-vs-ERROR 
classification | N/A | Low | Yes (@SEZ9) |
   
   # 5. Merge Recommendation
   
   ### Conclusion: Not recommended for merge
   
   **Blockers — must be fixed:**
   1. Issue 1 (Critical) — fix the Mockito/generics compile break in 
`SeaTunnelServerShutdownTest.java:55` (`doReturn(...).when(...)` or an explicit 
type witness on `getMap`), then get a fully green CI run on the resulting head 
before merge. This is the sole reason I'm not endorsing my own round-6 "Ready 
to merge after fixes" today — that conclusion was reached while CI was still 
queued, and the completed run tells a different story than I assumed it would.
   
   **Recommended fixes — non-blocking, all carried over and still valid:**
   - Issue 5 — add a round-trip test (or shared constant) tying the message 
format string to the regex; cheapest, highest-leverage of the remaining items.
   - Issue 3 — native-TTL `IMap.put` overload, non-destructive-until-processed 
marker consumption.
   - Issue 4 / Issue 6 — carry the graceful/non-graceful classification as an 
explicit typed flag rather than a payload-shape or free-text difference; would 
close both at once.
   - Issue 2 — correct the `onShutdown()` comment.
   - Issue 7 — a short docs note on the WARN/ERROR distinction for operators 
who alert on Zeta engine logs.
   
   **Overall assessment:** the design work across six rounds — including a 
second independent reviewer's real catch and a verified-correct fix for it — 
has been genuinely solid, and I stand by everything I said about the runtime 
logic in round 6: it fails safe (every ambiguous case defaults to the 
pre-existing `ERROR` behavior), and the mechanism itself is sound. But "ready 
to merge after fixes, contingent on CI" was not actually satisfied by this head 
— CI finished, and it's red for a real, deterministic reason. I should have 
circled back and confirmed the run before writing that conclusion instead of 
asking a maintainer to check it independently; I'm correcting that now rather 
than letting a stale "no blockers" stand on the PR. Once Issue 1 is fixed and 
CI is actually green on the resulting commit, I'd be comfortable with a 
maintainer merge decision — the remaining Medium/Low items are legitimate 
hardening work, not correctness gates.
   


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