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

   @SEZ9 Here's the grounding you asked for, checked directly against the 
current head (`b921264bb`, unchanged since my last review):
   
   **F4 — coordinator collection path, with the actual code:**
   
   `CoordinatorService.collectCdcEnumeratorProgress()` 
(`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/CoordinatorService.java:2190-2248`)
 is the implementation, and its own Javadoc (lines 2182-2188) states the 
contract directly: "Enumerator tasks follow normal coordinator-task placement 
and may execute on any member. Their current locations are derived from the 
coordinator-owned slot map... worker-local registration is not part of the 
ownership contract."
   
   Walking the method itself:
   - Lines 2195-2233: it iterates every running `JobMaster`'s physical plan, 
looks up each pipeline's `ownedSlotProfilesIMap` (the coordinator-owned slot 
assignment map), and for every coordinator vertex whose task group contains a 
`SourceSplitEnumeratorTask`, resolves that task group's current location from 
the slot map — this is the "derives enumerator task group locations from 
running job plans and coordinator-owned slot assignments" step.
   - Lines 2235-2247: it then calls the overloaded 
`collectCdcEnumeratorProgress(taskGroupsByWorker, collect, update, onFailure)`, 
where `collect` sends a 
`CollectCdcEnumeratorProgressOperation(taskGroupLocations)` to each resolved 
worker address via `NodeEngineUtil.sendOperationToMemberNode` — this is the 
"requests reports from the assigned members" step — and completions feed into 
`seaTunnelServer.getCdcProgressService().updateReports(reports)`.
   - The static overload at lines 2255-2274 fires all worker requests up front 
before attaching completion callbacks (matching the Javadoc note on that 
method: callbacks may run on Hazelcast completion threads and must stay 
non-blocking).
   
   This is a genuine request/poll model driven entirely by coordinator-side 
state (the slot map), with no worker-initiated registration step anywhere in 
the method. I also diffed this against `docs/en/developer/cdc-progress.md`'s 
"Runtime collection" section ("The active coordinator derives enumerator task 
group locations from running job plans and coordinator-owned slot assignments. 
It requests reports from the assigned members...") — the doc text and the code 
match phrase-for-phrase, not just in spirit. F4 is grounded and resolved.
   
   **F8 — explicit conclusion, since you asked me to state it:**
   
   Yes, the deep-immutability guarantee holds at every layer, verified against 
the actual (unchanged) production code, not just the new test:
   - `CdcProgressPosition`'s constructor does a real defensive copy — 
`Collections.unmodifiableMap(new LinkedHashMap<>(values))` — so mutating the 
caller's original map after construction cannot reach the stored copy, and the 
returned view itself rejects writes.
   - `CdcSnapshotSplitProgress` is `public final class` with four `final` 
fields and no setters — no mutation path exists post-construction regardless of 
what's passed in.
   - `CdcEnumeratorProgressReport`'s active-splits list is wrapped the same way 
(`Collections.unmodifiableList`), and the new 
`testEnumeratorReportActiveSplitWatermarksAreDeeplyImmutable` now exercises the 
full nested chain (`CdcEnumeratorProgressReport` → `CdcSnapshotSplitProgress` → 
`CdcProgressValue<CdcProgressPosition>` → the map), not just the outer list. I 
confirmed the test is discriminating (it would fail if the defensive copy were 
removed) rather than vacuously true.
   
   F8 is resolved with both source-level confirmation and direct regression 
coverage.
   
   **F7 pointer** (since @goutamadwant's ask was specifically for a method 
name): `CdcProgressModelTest.testEnumeratorReportRejectsInvalidCounts` 
(`seatunnel-api/src/test/java/org/apache/seatunnel/api/cdc/CdcProgressModelTest.java:155-180`)
 — asserts `IllegalArgumentException` for a negative snapshot-total count and 
for a completed-count that exceeds the total, both via the 
`CdcEnumeratorProgressReport` constructor directly.
   
   **Explicit record on all of F1-F8 against `b921264bb`:** all eight remain 
resolved from my side. Nothing in this head's one-file test diff touches F1-F6 
(those depend on production code that hasn't changed since my 2026-08-31 
line-by-line pass), and F7/F8 are now grounded above. I don't have an open item 
on this list.
   


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