DanielLeens commented on code in PR #11841:
URL: https://github.com/apache/seatunnel/pull/11841#discussion_r3919904823


##########
seatunnel-engine/seatunnel-engine-core/src/main/java/org/apache/seatunnel/engine/core/job/JobImmutableInformation.java:
##########
@@ -167,6 +168,14 @@ public boolean isRestoreJob() {
         return restoreMode != null && restoreMode.isRestore();
     }
 
+    public boolean isSavepointRestore() {

Review Comment:
   Good question to raise on a checkpoint/savepoint-adjacent field, but I don't 
think the premise holds once you trace where `JobImmutableInformation` actually 
gets serialized.
   
   `JobImmutableInformation` itself is never embedded in checkpoint/savepoint 
**state** storage (`PipelineState`/`CheckpointStorage`) — I grepped 
`seatunnel-engine-server/.../checkpoint/**` and there's no reference to it 
there. It's freshly *constructed* by the client at each resubmission, with 
`restoreMode` explicitly set to `SAVEPOINT`/`CHECKPOINT` by the current-version 
code doing the restoring 
(`ClientJobExecutionEnvironment`/`RestJobExecutionEnvironment`). So for the 
scenario you describe — user upgrades SeaTunnel, then restores from a savepoint 
taken earlier — the `JobImmutableInformation` used for that restore submission 
is built fresh by the new binary; it isn't deserialized from old bytes at all.
   
   The two places it genuinely round-trips through serialization are the 
RPC/REST submission wire and `runningJobInfoIMap` (used for master-node 
failover recovery within a live cluster — see the comment at 
`CoordinatorService.java:147`). Both go through `readData()` 
(`JobImmutableInformation.java:243-269`), which is **untouched by this PR** 
(verified via `git diff apache/dev...HEAD` on this file — the only changes are 
the two new accessors at lines 171-177 and a one-line Javadoc). That method's 
pre-existing fallback logic is:
   
   ```java
   restoreMode = isStartWithSavePoint ? RestoreMode.SAVEPOINT : 
RestoreMode.NONE;   // line 259
   restoreSourceJobId = isStartWithSavePoint ? jobId : null;                    
    // line 260
   if (hasRemainingBytes(in)) { ... }                                           
    // trailer, if present
   ```
   
   This derives `restoreMode` from the legacy boolean *before* checking for the 
trailer, so an old-format payload (pre-`RestoreMode`, pre-#11421) that was 
genuinely savepoint-restored has `isStartWithSavePoint=true` on the wire and 
correctly falls back to `RestoreMode.SAVEPOINT` — not "all false." 
`isSavepointRestore()`/`isCheckpointRestore()` are pure derivations of that 
already-correct `restoreMode` value, so they inherit the same backward-compat 
guarantee without adding any new deserialization logic themselves.
   
   So: no bug here that I can find — the compatibility path you're worried 
about was already handled before this PR touched the file, and this PR's own 
diff to this class is additive-only (two accessor methods + a comment). If you 
have a specific call path in mind where `JobImmutableInformation` *does* get 
persisted to durable storage and read back post-upgrade (rather than freshly 
built per submission), I'd genuinely like to see it — that would be a real 
pre-existing gap worth its own issue, just not one introduced by this diff.
   



##########
seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/master/JobMaster.java:
##########
@@ -284,7 +285,7 @@ public synchronized void init(long initializationTimestamp, 
boolean restart) thr
                                                             sink.getId()));
                                     JobMaster.handleSaveMode(
                                             ((SinkAction<?, ?, ?, ?>) 
sink).getSink(),
-                                            logicalDag.isStartWithSavePoint());

Review Comment:
   Traced this through `LogicalDag`/`LogicalDagGenerator` to answer precisely 
rather than just from the diff.
   
   `LogicalDag.isStartWithSavePoint` (`LogicalDag.java:63`) is **untouched** by 
this PR — it's still a plain `boolean` field. It's populated at DAG-generation 
time via `AbstractJobEnvironment.buildLogicalDagGenerator` → `new 
LogicalDagGenerator(actions, jobConfig, idGenerator, restoreMode.isRestore())` 
(`AbstractJobEnvironment.java:135`, also untouched), which is already the 
**broad** "any restore" truth value (true for both `SAVEPOINT` and 
`CHECKPOINT`), despite the field's savepoint-sounding name. So 
`logicalDag.isStartWithSavePoint()` and 
`jobImmutableInformation.isRestoreJob()`/`getRestoreMode().isRestore()` were 
already computing the identical boolean for every reachable path before this PR 
— the swap at `JobMaster.java:273` and `:287-288` is value-preserving, not a 
behavior change.
   
   On your actual question — should `RestoreMode` live on `LogicalDag` instead 
of `JobMaster` reaching into `JobImmutableInformation` directly: I'd argue the 
current shape is actually the *more* honest one, not the ambiguous one. 
`LogicalDag`'s field is a `boolean`, so if we wanted `LogicalDag` to 
distinguish `SAVEPOINT` vs `CHECKPOINT` for any future use, it would need to 
become a `RestoreMode`-typed field anyway — which would just be a second, 
derived copy of the exact same value `JobImmutableInformation` already carries 
authoritatively. `JobMaster` already holds a direct reference to 
`jobImmutableInformation` (it's a constructor field, see 
`jobImmutableInformation.getJobConfig()` two lines above at 
`JobMaster.java:242`) and is the object that originally caused `logicalDag`'s 
own flag to be set in the first place (client-side, via the same 
`restoreMode`). So going straight to `jobImmutableInformation.getRestoreMode()` 
at the master-side call site removes a layer of indirection 
 through a misleadingly-named boolean proxy, rather than introducing a new 
coupling — `JobMaster` isn't reaching into someone else's internals here, it's 
using the source of truth it already owns.
   
   I don't see a separation-of-concerns issue in practice, but I take the point 
that `LogicalDag.isStartWithSavePoint`'s name is now doubly misleading (it was 
already broad-but-narrowly-named before this PR, and this PR's own 
`CheckpointCoordinator` parameter rename — `isStartWithSavePoint` → 
`isRestoreJob` — fixes exactly this class of naming problem elsewhere in the 
same diff). Renaming 
`LogicalDag.isStartWithSavePoint()`/`setStartWithSavePoint()` to something like 
`isRestoreJob()` would be a reasonable, low-risk follow-up for consistency, 
though it's a wire-serialized field (`LogicalDag.java:158,174`) so it'd need 
the same "keep the old accessor name, just clarify semantics via Javadoc" 
treatment this PR already used on 
`JobImmutableInformation.isStartWithSavePoint()` rather than an outright 
rename, to avoid a gratuitous compat-sensitive touch. Happy to see that filed 
separately if you'd like to pursue it — not a blocker for this PR either way.
   



##########
seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java:
##########
@@ -1394,9 +1396,12 @@ public PassiveCompletableFuture<Void> submitJob(
                     try {
                         JobImmutableInformation 
submittedJobImmutableInformation =
                                 
deserializeJobImmutableInformation(jobImmutableInformation);
+                        boolean isSavepointRestore =

Review Comment:
   This is a fair, sharp catch — I checked and you're right that it's dead 
inside this method: `submitJob(long jobId, Data jobImmutableInformation, 
boolean isStartWithSavePoint)` (`CoordinatorService.java:1374-1375`) never 
references the `isStartWithSavePoint` parameter anywhere in its ~90-line body 
once you get past the signature — I grepped the full method and the only 
occurrence is the declaration itself. Everything downstream uses 
`submittedJobImmutableInformation.isSavepointRestore()`, derived from the 
deserialized payload (line 1399), exactly as you traced.
   
   Two things worth putting on the record here, though:
   
   1. **This exact tension was already raised and resolved in this PR's own 
thread**, before your comments — see @SEZ9's question and my answer in the 
top-level conversation on 2026-08-20/21 (issue comments). The explicit, 
deliberate decision was: keep the legacy wire argument as-is in this PR for 
compatibility, and treat any deprecation/signature cleanup of the now-unused 
parameter as a dedicated follow-up once every reader has migrated off it, 
rather than bundling a signature change into an already-reviewed, additive 
refactor. That's consistent with what you're describing — the parameter 
*should* eventually go away or be replaced with an explicit `RestoreMode`, just 
not in this diff.
   
   2. **It's not purely inert further downstream, which is why I wouldn't call 
it "obsolete" outright** — the same boolean is also what gets written onto 
`SubmitJobOperation`'s own wire format (`writeInternal`/`readInternal`, 
untouched by this PR — see your other comment on `BaseService.java:1360`, I've 
answered the wire-level angle there). That's a genuinely separate concern from 
`CoordinatorService.submitJob`'s *internal* decision-making: during a rolling 
upgrade, an old-version master node doesn't know about 
`JobImmutableInformation.isSavepointRestore()`/`RestoreMode` at all and needs a 
coherent `isStartWithSavePoint` value to still be present on the wire it can 
read. So the parameter has to keep existing at the RPC boundary even after it 
becomes irrelevant to this method's own logic — which is exactly the "kept for 
wire compat, no longer trusted for decisions" framing in my Compatibility 
Impact section of the review above.
   
   So: agreed this deserves a real cleanup, agreed it isn't currently reflected 
in the method signature, but I'd treat it as the deprecation follow-up already 
discussed with @SEZ9 rather than something this PR needs to absorb — happy to 
see a tracking issue opened for "thread `RestoreMode` explicitly through 
`CoordinatorService.submitJob`/`SubmitJobOperation` and retire the legacy 
boolean" if you or @JeremyXin want to file it.
   



##########
seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/BaseService.java:
##########
@@ -1360,7 +1360,7 @@ protected JsonObject submitJobInternal(
                             new SubmitJobOperation(

Review Comment:
   Same underlying observation as your comment on 
`CoordinatorService.java:1399`, traced through the Hazelcast operation layer 
instead — I answered the "should this be threaded explicitly / is it now dead" 
question there. On the "silent override" framing specifically, though, I don't 
think there's an actual override risk today, only redundancy:
   
   At both of `BaseService.java`'s own producers of this value — the cross-node 
RPC path at line 1360 you quoted, and the same-node direct path at line 1397 
(`submitJob(...)` private helper) — the wire boolean passed into 
`SubmitJobOperation`/`CoordinatorService.submitJob` is 
`jobImmutableInformation.isSavepointRestore()`. That's the *same* 
`RestoreMode`-derived value that `CoordinatorService.submitJob` later 
re-derives from the deserialized payload 
(`submittedJobImmutableInformation.isSavepointRestore()`, 
`CoordinatorService.java:1399`). Same is true of `ClientJobProxy.java:77`, the 
third producer. So for every currently-reachable caller in this codebase, the 
value written onto the wire and the value the coordinator re-derives are 
byte-for-byte identical by construction — there's no case in the current diff 
where the outer parameter says one thing and the payload says another.
   
   The scenario where they *could* genuinely diverge is a caller that predates 
this consolidation (a not-yet-migrated or old-version producer sending a 
stale/inconsistent `isStartWithSavePoint` alongside a payload with a different 
real `RestoreMode`) — and that's precisely what the new 
`CoordinatorServiceJobCleanupTest.testSubmitCheckpointRestoreDoesNotTrustLegacySavepointFlag`
 test simulates and asserts against (submits `RestoreMode.CHECKPOINT` with the 
wire boolean forced to `true`, asserts the checkpoint path — not savepoint's 
cleanup-ownership branch — is what runs). So `CoordinatorService` deliberately 
no longer trusts the RPC-level boolean for exactly the reason you're pointing 
at: it *could* become stale/inconsistent from a future or mixed-version caller, 
even though every current in-repo caller keeps it consistent. That's a 
hardening, not a bug being papered over.
   
   `SubmitJobOperation.writeInternal`/`readInternal` (the code you pasted) is 
untouched by this PR and needs to keep writing/reading that boolean regardless, 
for rolling-upgrade wire compatibility with an old-version node on the other 
end — see my fuller answer on the `CoordinatorService.java:1399` thread for why 
it can't just be dropped from the wire format even though it's no longer 
trusted for decisions on the receiving side.
   



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