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

   Hi @SEZ9, thank you for coming back to this thread, and for confirming that 
`dc2c64103c` closes the compile break. This is my own PR, so I am answering as 
the author and have re-checked everything below on the current head 
`dc2c64103c` (unchanged since my last comment) and in the fork CI run for that 
head.
   
   ## Direct answers
   
   **F1 - where the listener is registered, where the write happens, and 
whether a real lifecycle path is tested**
   
   - Registration: `SeaTunnelServer.init()` calls 
`getLifecycleService().addLifecycleListener(this)` at 
`SeaTunnelServer.java:154` (the class implements `LifecycleListener`).
   - Dispatch: `stateChanged(...)` at `SeaTunnelServer.java:239-245` sends 
`SHUTTING_DOWN` to `markLocalGracefulMemberRemoval()` (`:284-286`) and 
`STARTED` to the clear (`:292-294`).
   - The IMap write is `gracefulMemberRemovalIMap.put(thisAddress, 
System.currentTimeMillis(), TTL, MILLISECONDS)` at 
`SeaTunnelServer.java:309-313`, inside `updateLocalGracefulMemberRemovalMarker` 
(`:300-326`). `ManagedService.shutdown()` (`:247-278`) no longer touches the 
map, and 
`SeaTunnelServerShutdownTest.shouldNotMarkGracefulMemberRemovalDuringManagedServiceShutdown`
 pins that.
   - Real lifecycle test: 
`ClusterFaultToleranceIT.testGracefulShutdownPublishesMemberRemovalMarker` 
(`ClusterFaultToleranceIT.java:79-146`) starts two real members with 
`SeaTunnelServerStarter.createHazelcastInstance`, registers an entry listener 
on node2's real map proxy (`:108-119`), calls `node1.shutdown()` (`:121`), and 
asserts the marker arrives and is accepted by 
`CoordinatorService.isGracefulMemberRemovalMarkerValid` (`:124-131`). The 
mock-based `SeaTunnelServerShutdownTest` only covers the branch logic.
   - Evidence that it ran: in fork run `35240527669` (this head, attempt 3), 
the `engine-v2-it (8, ubuntu-latest)` log reports `Tests run: 10, Failures: 0, 
Errors: 0, Skipped: 2 - in ClusterFaultToleranceIT`. The two skips match the 
two `@Disabled` tests that already exist in that file on `dev` and on this 
head. That the new test itself passed is an inference from the class-level 
result; I did not find a per-method line for it.
   - What that test does not cover: it exercises 
`HazelcastInstanceImpl.shutdown()`, not the JVM shutdown-hook thread (SIGTERM) 
path. That the hook reaches the same `LifecycleServiceImpl.shutdown` and fires 
`SHUTTING_DOWN` is from my reading of the Hazelcast 5.1 sources, not from a 
test in this PR.
   
   **F4 - TTL versus timestamp**
   
   Both, with different jobs. Eviction is a Hazelcast-side per-entry TTL: the 
4-argument `put(key, value, ttl, unit)` at `SeaTunnelServer.java:309-313` with 
`Constant.GRACEFUL_MEMBER_REMOVAL_MARK_TTL_MILLIS` = 5 minutes 
(`Constant.java:72`). Separately, 
`CoordinatorService.isGracefulMemberRemovalMarkerValid` 
(`CoordinatorService.java:2157-2161`) checks `abs(now - markedAt) <= 5 min`. 
That second check does compare the departing member's 
`System.currentTimeMillis()` with the coordinator's clock, so it tolerates up 
to 5 minutes of skew; beyond that the marker is rejected and the failure stays 
at `ERROR`, which is the safe direction. So the "no eviction" part of F4 is 
resolved by the native TTL, and the timestamp is only a secondary guard.
   
   **F2 - destructive removal and failover**
   
   The read is non-destructive: `gracefulMemberRemovalIMap.get(lostAddress)` at 
`CoordinatorService.java:2185`. The clear runs only after the propagation loop 
(`:2222-2238`), and only when no master-switch restore is in flight 
(`canClearGracefulMemberRemovalMarker`, `:2169-2172`, called at `:2239-2245`). 
It is a value-conditional `remove(lostAddress, markedAt)` (`:2207`), so it 
cannot delete a newer marker from a restarted member at the same address. When 
the clear is skipped, or a coordinator dies mid-processing, the entry stays in 
the IMap until its TTL, and the new master classifies through 
`PhysicalVertex.checkTaskGroupIsExecuting` (`PhysicalVertex.java:239-250`), 
which reads the same marker without clearing it (`:278-288`) and applies the 
same helpers (`recordMemberRemovedFailure`, `:760-770`). Unit tests: 
`shouldRetainGracefulMemberRemovalMarkerDuringMasterSwitchRecovery` and 
`shouldClassifyMissingWorkerDuringMasterFailoverAsGraceful`. One nit I noticed 
while re-reading: the Ja
 vadoc at `PhysicalVertex.java:273-277` still says "one-time" marker, which is 
stale wording now that the read is non-destructive.
   
   ## Per-finding status on `dc2c64103c`
   
   | Finding | Status | Evidence |
   |---|---|---|
   | F1 marker written while PASSIVE | Fixed in source; publication exercised 
by a real two-node IT that passed on this head | `SeaTunnelServer.java:154`, 
`:239-245`, `:309-313`; `ClusterFaultToleranceIT.java:79-146` |
   | F2 consumed before processing / failover | Fixed | 
`CoordinatorService.java:2185`, `:2207`, `:2239-2245`; 
`PhysicalVertex.java:278-288` |
   | F3/F6 `JobException` replaced by a string | Fixed; payload shape unchanged 
| `buildMemberRemovedFailureState`, `CoordinatorService.java:2144-2149`; 
`shouldKeepThrowablePayloadForMemberRemovedFailureState` |
   | F4 eviction | Fixed by native TTL; timestamp is a secondary guard | 
`SeaTunnelServer.java:309-313`; `CoordinatorService.java:2157-2161` |
   | F5/F8 message matching and duplicated template | Fixed | no `Pattern` or 
`.matches(` in `PhysicalVertex.java`; single producer 
`buildMemberRemovedOfflineMessage` (`CoordinatorService.java:2133-2137`), 
called from `:2146` and `PhysicalVertex.java:768`; the classification is a 
typed boolean |
   | F7 docs | Present | `docs/en/engines/zeta/state-storage-and-recovery.md` 
(new section at line 155) and the zh counterpart |
   
   ## What CI actually shows on this head
   
   I read fork run `35240527669` (attempt 3, finished 2026-09-18T14:13Z; the 
apache-side `Build` check finished as failure at 14:14Z). It is red, so I 
cannot call this PR verified or ready to merge yet.
   
   - The compile break is fixed: `unit-test` passed on Java 8 and 11 on both 
ubuntu and windows, and `engine-v2-it (11)` passed. 79 of 94 jobs succeeded, 11 
were skipped, 3 failed, and 1 was cancelled.
   - `engine-v2-it (8)` failed with 204 tests run, 1 failure and 1 error: 
`BackpressureSlowSinkIT.testCheckpointsKeepCompletingUnderSustainedBackpressure`
 (only 1 of the expected 3 additional checkpoints), and 
`SplitClusterFaultToleranceIT.testStreamJobCancelResolvesWhenWorkerCrashesBeforeCancelAck`
 (`expected: <CANCELED> but was: <FAILED>`). Both match the targets of open PRs 
#12316 and #12311, so a `dev` sync would not clear them. My diff does not touch 
the vertex-state filter or the FAILED target in `makeTasksFailed` 
(`CoordinatorService.java:2256-2265`); it only adds the classification 
argument. That these two are unrelated to my change is my inference, not 
something a run has proven.
   - `all-connectors-it-2 (8)` and `(11)` both failed on 
`OpengaussCDCIT.testAddFieldWithRestore` (`ConditionTimeout` at `:476`). This 
is a CDC connector test unrelated to the engine shutdown path, as far as I can 
tell.
   - `jdbc-connectors-it-part-1 (11)` was cancelled at 14:13Z after running 
about two hours, with no failed test in its log. I have not established why it 
was cancelled.
   
   ## Runtime evidence for the WARN downgrade (the acceptance list from my last 
comment)
   
   Counts below are from the `engine-v2-it (8)` log only, which is one job of 
the run.
   
   - `... end with state FAILED due to node offline` at `WARN`: 125 occurrences 
(baseline on the previous head: 0). No `ERROR`-level offline line remains in 
that log.
   - `NoDataMemberInClusterException`: 0 occurrences.
   - `Failed to clear graceful member removal marker`: 13, all 
`HazelcastInstanceNotActiveException`, logged from the coordinator's clear 
(`CoordinatorService.java:2207`). This is not fully gone; it fires when the 
coordinator itself is stopping.
   - `Failed to mark graceful member removal marker`: 39, all 
`HazelcastInstanceNotActiveException`, so the criterion "gone for shutdown() 
paths" is not literally met. I traced four of them in 
`ClusterFaultToleranceIT`, plus one in 
`SplitClusterPendingJobLifecycleFailoverIT`, and they are calls to `shutdown()` 
on an instance that was already down: for example 
`testBatchJobRestoreIn2NodeWorkerDown` shuts node2 down in the body (`:441`) 
and again in `finally` (`:470`), and the same pattern holds at `:540`/`:591`, 
`:660`/`:690` and `:763`/`:809`. `LifecycleServiceImpl.shutdown` fires 
`SHUTTING_DOWN` on every call, so the listener tries a write against a dead 
instance. I have not traced the other 34 individually, so I cannot claim they 
are all of this kind.
   
   ## Remaining gaps, stated plainly
   
   1. CI on this head is red (list above). I am not treating the PR as ready 
until the failed jobs are re-run and pass, or are shown to be unrelated.
   2. No test asserts the `WARN` versus `ERROR` classification in a real 
cluster. The engine-level evidence is the log counts above; the classification 
rules themselves are covered by unit tests only (`PhysicalVertexTest`, 
`CoordinatorServiceMemberRemovedTest`).
   3. My earlier recommendation N2 is still open on this head: the marker `put` 
(`SeaTunnelServer.java:309`) and the `STARTED` clear (`:315`) are synchronous 
calls on lifecycle threads. Together with the stack-trace noise above, a 
bounded async call, or skipping the write when the instance is already not 
running, would be a reasonable follow-up.
   4. N4 (one sentence in the docs about the 5-minute TTL versus 
`hazelcast.max.no.heartbeat.seconds`) is also still open.
   
   I will do a final pass on the next green run of this head. Thanks again for 
the careful review.
   


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