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

   Re-reviewed from scratch at head `20538ae1ebe78cc2d35871fbd02e1d79b960e154` 
— this is the **same commit** I reviewed on 2026-08-16 (no new commits landed 
on this branch in the last 3 days). Rather than re-post a copy of that review, 
I re-traced the whole call chain independently again (own-code trust boundary 
applies to me too), re-verified every claim in the 2026-08-16 comment against 
the current source, and I'm folding the result into one review below so there's 
a single up-to-date place to look. Where a prior finding still holds I say so 
and cite what I re-checked; I did not find anything that changes the earlier 
conclusions.
   
   # What Problem Does This PR Solve?
   - User pain point: Zeta's `/overview` endpoint only gives cluster-wide slot 
totals; there is no way to ask "how many slots is running job X holding right 
now, and on which workers/pipelines". Operators diagnosing a stuck or 
over-provisioned cluster have to reconstruct this from logs.
   - Fix approach: adds a read-only REST endpoint, `/running-jobs/slot-usage` 
(v2, Jetty) and `/hazelcast/rest/maps/running-jobs/slot-usage` (v1), backed by 
a master-side operation and a stateless aggregation helper 
(`RunningJobSlotUsageBuilder`) that intersects the coordinator's owned-slot map 
with the resource manager's assigned-slot snapshot, grouped by job / pipeline / 
worker.
   - One-sentence summary: a purely additive, read-only per-running-job 
slot-usage diagnostics endpoint for the Zeta REST API.
   
   # 1. Code Change Review
   
   ## 1.1 Core Logic Analysis
   
   Runtime path, re-traced against the current head:
   
   ```text
   HTTP GET /running-jobs/slot-usage                       (Jetty v2)
     -> JettyService.java:190,230  servlet mapped via convertUrlToPath() -> 
"/running-jobs/slot-usage/*"
     -> RunningJobSlotUsageServlet.doGet()                  
[RunningJobSlotUsageServlet.java:41-45]
     -> RunningJobSlotUsageService.getRunningJobSlotUsageJson()  
[RunningJobSlotUsageService.java:35-43]
          |
          +-- this node IS master (getSeaTunnelServer(true) != null)
          |     -> RunningJobSlotUsageBuilder.build(server)      [in-process, 
on the Jetty thread]
          |
          +-- this node is NOT master
                -> 
NodeEngineUtil.sendOperationToMasterNode(GetRunningJobSlotUsageOperation)
                     [NodeEngineUtil.java:33-41, service name = 
SeaTunnelServer.SERVICE_NAME]
                -> master: GetRunningJobSlotUsageOperation.run()  
[GetRunningJobSlotUsageOperation.java:92-96]
                     runs on a Hazelcast generic operation thread (see Issue 3)
                -> RunningJobSlotUsageBuilder.build(server)
   
   HTTP GET /hazelcast/rest/maps/running-jobs/slot-usage    (v1)
     -> RestHttpGetCommandProcessor.handle() ordering:
          RUNNING_JOBS_SUMMARY -> RUNNING_JOBS_SLOT_USAGE -> RUNNING_JOBS   
[RestHttpGetCommandProcessor.java:161-166]
     -> handleRunningJobsSlotUsage() -> same RunningJobSlotUsageService
   
   Aggregation (RunningJobSlotUsageBuilder.build(server), both entry points 
converge here)
     step 1  IMAP_RUNNING_JOB_INFO.keySet().forEach(...)          
[RunningJobSlotUsageBuilder.java:59-66]
             filter: 
server.getCoordinatorService().shouldShowAsRunningJob(jobId)
               -> CoordinatorService.java:1502 -> getJobStatus(jobId) -> 
!status.isEndState()
     step 2  server.getCoordinatorService().getInitializedResourceManager()  
[line 68-69]
               -> plain volatile field read, no lazy init 
(CoordinatorService.java:1247-1249)
             resourceManager == null ? emptyList() : 
getAssignedSlots(emptyMap())  [line 70-73]
               -> AbstractResourceManager.java:328-332, streams 
registerWorker.values()
                  -> WorkerProfile.getAssignedSlots() (master-local, 
heartbeat-refreshed state)
     step 3  ownedSlotProfiles.forEach((pipelineLocation, taskGroupSlots) -> 
...)  [line 88-92, 414-433]
             per pipeline: keep only slots whose SlotKey is present in the 
assigned-slot set
               SlotKey = (worker, slotId, ownerJobId, sequence)   [line 450-490]
     step 4  SlotUsage.addSlot() -> slotCount / pipelineSlotCounts / 
workerSlotCounts
             [line 503-509] -> toResponse() [line 511-518] -> 
JsonUtils.toJsonString
   ```
   
   Things I checked myself rather than take on faith, because each one would 
have been a real bug if wrong:
   
   - `SlotProfile.ownerJobID` is written/read in `writeData`/`readData` — 
confirmed via `grep` on the constructor/serialization methods. Since `SlotKey` 
includes `ownerJobId`, a non-serialized owner id would make the intersection 
miss on every cross-node call and `slotCount` would be permanently `0` on any 
real (multi-node) cluster. It's on the wire, so this genuinely works.
   - `PipelineLocation` carries `jobId` as one of its two fields 
(`PipelineLocation.java`), so `usageByJob.get(pipelineLocation.getJobId())` at 
line 423 can never attribute one job's owned-slot entries to a different job's 
`SlotUsage`. No cross-job leakage from the primary grouping key.
   - New serializer class id `17` (`GET_RUNNING_JOB_SLOT_USAGE_OPERATION`) does 
not collide — `ClientToServerOperationDataSerializerHook` on `dev` tops out at 
`16` (`GET_JOB_TASK_MAPPING_OPERATION`); this PR is insert-only there.
   - Servlet routing is correct on both API versions: `convertUrlToPath()` 
appends `/*`, so `/running-jobs/slot-usage/*` is a strictly longer path than 
`/running-jobs/*` and wins under the servlet longest-prefix rule; and the v1 
`startsWith` chain places `REST_URL_RUNNING_JOBS_SLOT_USAGE` before the more 
general `REST_URL_RUNNING_JOBS`, so it's not swallowed by the existing branch — 
same technique already used for `REST_URL_RUNNING_JOBS_SUMMARY`.
   - `getService()` resolves without `GetRunningJobSlotUsageOperation` 
overriding `getServiceName()`, because 
`NodeEngineUtil.sendOperationToMasterNode` builds the invocation with 
`SeaTunnelServer.SERVICE_NAME` explicitly — same pattern as the other no-arg 
operations in this package (e.g. `GetNodeHttpPortOperation`).
   - Read-only, confirmed by walking every call the new code makes: 
`IMap.keySet()`, `IMap.forEach`, `shouldShowAsRunningJob` (reads only), 
`getInitializedResourceManager()` (field read), 
`AbstractResourceManager.getAssignedSlots` (stream over `registerWorker`). 
Nothing calls `requestResource`, `releaseResource`, `slotActiveCheck` or 
`heartbeat`. No slot allocation/deallocation side effect from querying this 
endpoint.
   - `getInitializedResourceManager()` (not the lazy `getResourceManager()`) is 
the right call for a read path — the sibling `GetOverviewOperation.java:82` 
still uses the lazy-init variant, so `/overview` can trigger cluster RPCs and 
runtime-state creation just from being polled. This endpoint deliberately 
avoids that. Good instinct; see Issue 1 for the cost of that choice.
   
   What "usage" actually measures, verified against the code rather than the 
field names: the numerator is the set of `IMAP_OWNED_SLOT_PROFILES` entries 
(coordinator-side, cluster-replicated, written at deploy time) that also appear 
in `ResourceManager.getAssignedSlots(...)` (master-local, heartbeat-refreshed 
`registerWorker` snapshot). That's "coordinator-owned task-group slots the 
current master also currently believes are assigned" — an intersection of two 
views with different freshness characteristics, not a single authoritative 
read. The docs phrase this accurately ("assigned task group slots owned by the 
running job"), so docs and code agree; the residual gap is between the code and 
the endpoint's own name/field name (see Issue 7).
   
   ## 1.2 Compatibility Impact — **Fully compatible**
   
   - Purely additive: new constant, new servlet, new service, new operation, 
new serializer class id. The diff on `JettyService`, 
`RestHttpGetCommandProcessor`, `RestConstant`, 
`ClientToServerOperationDataSerializerHook` is insert-only (13 files, 626 
additions, 0 deletions — confirmed via `git diff dev...HEAD --stat`).
   - No existing endpoint's response shape changed (`/running-jobs`, 
`/running-jobs/summary`, `/overview`, `/job-info` untouched).
   - No config option added, renamed, or re-defaulted.
   - No checkpoint/savepoint/state-serialization/job-restore surface touched.
   - Rolling-upgrade nuance worth one line in the PR description: operation 
class id `17` only ever travels worker -> master. In a mixed-version cluster 
where a new-version node serves this REST call but an old-version node is 
master, the old master's `ClientToServerOperationDataSerializerHook` factory 
throws `IllegalArgumentException("Unknown type id 17")`. That failure is 
confined to this new endpoint (no old client calls it), so it doesn't break any 
existing upgrade path, but it's worth documenting.
   
   ## 1.3 Performance / Side-Effect Analysis
   
   - Iteration safety re-checked: `IMap.forEach` walks a materialized 
`entrySet()` snapshot, not a live view; 
`AbstractResourceManager.registerWorker` is a `ConcurrentHashMap` with a 
weakly-consistent iterator. `SlotUsage` instances are confined to one call, no 
shared mutable state, no locks. No `ConcurrentModificationException` is 
possible.
   - Torn reads are real and inherent to the design: the running-job id set, 
the assigned-slot snapshot, and the owned-slot map are read at three different 
instants with no consistency barrier between them. A job that completes 
mid-call can show a stale count; a job that starts mid-call can show a partial 
one. Acceptable for a diagnostics endpoint, but combined with Issue 1 below, 
the caller can't tell a torn read from a real zero.
   - Cost re-verified: `IMAP_OWNED_SLOT_PROFILES.forEach` (line 88-92) pulls 
and deserializes **every** pipeline entry in the cluster to the calling member, 
including entries for jobs immediately discarded at line 423-426 because 
they're not in `runningJobIds`. `IMAP_RUNNING_JOB_INFO.keySet()` is likewise a 
distributed call. So per-request wire cost scales with total cluster state, not 
with the number of running jobs, on every single request — no cache, no TTL, no 
rate limit. To be fair to the PR, this same "read everything, filter after" 
shape already exists elsewhere in the engine for periodic/internal paths; the 
difference here is this one is externally, arbitrarily triggerable.
   - The bigger cost isn't CPU, it's thread blocking, covered as Issues 2 and 3 
below.
   
   ## 1.4 Error Handling and Logging
   
   - No per-job 404 path exists by design — the endpoint always returns the 
full running-job list, and a job that ends mid-call just comes back with a 
stale/zero count rather than throwing. That's a reasonable shape for a 
bulk-listing diagnostic endpoint.
   - Re-verified: the whole new path (servlet, service, operation, builder) 
contains zero log statements. Combined with Issue 1, the silent-zero condition 
leaves no trace anywhere.
   - Re-verified: `RunningJobSlotUsageService.getRunningJobSlotUsageJson()` 
calls `getSeaTunnelServer(true)` and, on the forwarding branch, `.join()`s the 
master invocation with no try/catch (`RunningJobSlotUsageService.java:37-41`). 
If mastership is lost between that check and the local `build()` call on the 
other branch, `getCoordinatorService()` throws `SeaTunnelEngineException`, and 
that reaches Jetty as an unhandled 500. I checked `BaseServlet.doGet` 
implementations for a sibling (`OverviewServlet`) and confirmed this exact gap 
already exists there too — this is not a regression introduced by this PR, it's 
the established (imperfect) pattern in this REST layer. I'm flagging it as Low 
for that reason, matching my own 2026-08-16 assessment.
   
   ---
   
   **Issue 1: `slotCount: 0` is both the "no slots yet" answer and the 
"couldn't read the data" answer**
   - **Location:** 
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/trace/RunningJobSlotUsageBuilder.java:68-73`,
 `:446-448`
   - **Problem description:** `getInitializedResourceManager()` never 
lazy-inits by design, so it's `null` whenever the resource manager hasn't been 
touched yet on this master (fresh master after failover, or before the first 
slot request). When that happens, `assignedSlots` is forced to 
`Collections.emptyList()`, the assigned-slot set is empty, `isAssignedToJob` 
rejects every entry, and every running job renders with `slotCount: 0` — even 
though `IMAP_OWNED_SLOT_PROFILES` may already hold full ownership data for 
those jobs. The docs (`docs/en/engines/zeta/rest-api-v2.md`) explicitly make 
`0` a legitimate value ("a running job that has not received any slot yet is 
returned with `slotCount` set to `0`"), so the response has no way to 
distinguish "genuinely zero" from "couldn't verify".
   - **Potential risk:** the moment an operator is most likely to query this 
endpoint — right after a master failover, while diagnosing whether jobs kept 
their resources — is exactly the moment it's most likely to lie and say every 
job has zero slots, with nothing logged to contradict it.
   - **Best improvement:** make unavailability explicit instead of encoding it 
as `0`. Option A: add a `slotSourceAvailable`/`degraded` boolean to each 
response entry when `resourceManager == null`, and document it. Option B: when 
the resource manager isn't initialized yet but `ownedSlotProfiles` is 
non-empty, fall back to counting directly from `IMAP_OWNED_SLOT_PROFILES` and 
mark that response as coordinator-only/unverified. Either way, log once at WARN 
with the affected job count. Add a unit test for this exact case (resource 
manager `null`, `ownedSlotProfiles` non-empty).
   - **Severity:** Medium
   
   **Issue 2: `getCoordinatorService()` runs once per running job inside the 
filter lambda, and it can sleep**
   - **Location:** 
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/trace/RunningJobSlotUsageBuilder.java:59-66`
   - **Problem description:** re-verified line by line — 
`runningJobInfo.keySet().forEach(jobId -> { if 
(server.getCoordinatorService().shouldShowAsRunningJob(jobId)) ... })` calls 
`getCoordinatorService()` once per entry in `IMAP_RUNNING_JOB_INFO`. That's not 
a cheap accessor: `SeaTunnelServer.getCoordinatorService()` 
(`SeaTunnelServer.java:286-303`) runs up to 3 retries of `Thread.sleep(500)` 
plus a WARNING log per iteration while the coordinator isn't active yet. With N 
running jobs and a coordinator that's briefly inactive (e.g. mid-transition), a 
single REST call can turn into up to N x 1.5s of sleeping plus N log lines. The 
reference is loop-invariant — trivial to hoist.
   - **Potential risk:** on a cluster with many running jobs during a master 
transition, one REST call blocks its thread for a long time and floods the log; 
combined with Issue 3, that thread is a shared Hazelcast operation thread, not 
a request-scoped one.
   - **Best improvement:** resolve `CoordinatorService coordinatorService = 
server.getCoordinatorService();` once before the loop and reuse the local 
reference for both the `shouldShowAsRunningJob` filter and the 
`getInitializedResourceManager()` call. Two-line fix, removes the O(N) 
amplification entirely.
   - **Severity:** Medium
   
   **Issue 3: The master operation aggregates inline on a shared Hazelcast 
operation thread; the directly comparable sibling operation deliberately 
doesn't**
   - **Location:** 
`seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/operation/GetRunningJobSlotUsageOperation.java:92-96`
   - **Problem description:** `run()` calls 
`RunningJobSlotUsageBuilder.build(service)` directly on the invoking thread. 
The operation declares no partition id, so it runs on a small, shared, 
cluster-wide generic operation thread pool. I compared this against 
`GetRunningJobMetricsOperation` (same package) and confirmed it deliberately 
offloads its body via `CompletableFuture.supplyAsync(..., 
getNodeEngine().getExecutionService().getExecutor("get_running_job_metrics_operation"))`
 — a dedicated named executor — precisely to keep this kind of work off the 
shared generic pool (the calling thread still blocks on `future.get()`, but a 
saturated dedicated executor only stalls this operation type, not every 
unrelated cluster operation sharing the generic pool). 
`GetRunningJobSlotUsageOperation` also declares `AllowedDuringPassiveState`, 
i.e. it can run while the cluster is passive — exactly when the coordinator is 
most likely to be inactive and the Issue 2 sleep path most likely to fir
 e.
   - **Potential risk:** blocking generic operation threads degrades unrelated 
cluster operations, not just this endpoint. A polling UI/monitor plus a 
coordinator transition is a plausible route to operation-thread pressure across 
the cluster, which on a production Zeta cluster is a job-availability concern, 
not just a monitoring inconvenience.
   - **Best improvement:** mirror `GetRunningJobMetricsOperation` — offload 
`build(...)` to a named executor from `getNodeEngine().getExecutionService()`. 
Fixing Issue 2 reduces the blocking duration but doesn't remove the 
distributed-map read cost from the shared operation thread, so both fixes are 
needed together.
   - **Severity:** Medium
   
   **Issue 4: Every request materializes the entire cluster-wide owned-slot 
map, including entries it immediately discards**
   - **Location:** `RunningJobSlotUsageBuilder.java:88-92`, `:423-426`
   - **Problem description:** `ownedSlotProfiles.forEach(...)` on a Hazelcast 
`IMap` resolves to `Map.forEach` over `entrySet()` — a distributed read that 
pulls and deserializes every pipeline's `Map<TaskGroupLocation, SlotProfile>` 
in the cluster to the calling member, before line 423-426 throws away 
everything not belonging to a running job. No caching, no TTL, no rate limit on 
the endpoint itself.
   - **Potential risk:** on a cluster with many pipelines, a tight polling loop 
turns into repeated full-map deserialization on the master, made worse by Issue 
3 since it happens on a shared operation thread.
   - **Best improvement:** narrow the read to running-job keys first (e.g. a 
predicate-based `IMap` read, or fetch by known `PipelineLocation` keys once 
`runningJobIds` is known) instead of a full `forEach`; alternatively a 
short-TTL cache in front of the aggregation, following the existing running-job 
DAG cache pattern already used elsewhere in `BaseService`.
   - **Severity:** Low
   
   **Issue 5: No logging anywhere in the new path; mastership loss on the 
forwarding branch surfaces as an unhandled 500**
   - **Location:** `RunningJobSlotUsageService.java:35-43`; 
`RunningJobSlotUsageBuilder.java:48-75`
   - **Problem description:** covered under 1.4 above — zero log statements in 
the whole feature, and an un-caught `.join()` on the forwarding path. Confirmed 
this matches the pre-existing pattern in sibling servlets, so it's a 
consistency nit rather than a regression, which is why I keep it Low rather 
than blocking on it.
   - **Potential risk:** a monitoring scrape sees an opaque 500 during a 
routine master transition instead of a clean transient error, and the Issue 1 
silent-zero condition is undiagnosable after the fact because nothing was 
logged.
   - **Best improvement:** capture the server reference once instead of two 
separate `getSeaTunnelServer` calls; wrap the forwarded `.join()` in a 
try/catch that logs at WARN; add one WARN when 
`getInitializedResourceManager()` returns `null` while running jobs exist.
   - **Severity:** Low
   
   **Issue 6: Test coverage doesn't reach the degraded path, the public entry 
point, or the REST layer**
   - **Location:** 
`seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/trace/RunningJobSlotUsageBuilderTest.java`
   - **Problem description:** the one test that exists is genuinely good: it 
exercises the package-private 3-arg overload, asserts real numeric values 
rather than a bare "200 OK", and correctly covers the "slot reused by a 
different job" race by giving the reused slot a different `ownerJobId` in the 
`assignedSlots` list than the one recorded in `ownedSlotProfiles` — I confirmed 
the assertion (`firstPipelineSlotCounts.containsKey(2)` is false) genuinely 
exercises `SlotKey` rejecting a stale/reassigned entry, which is the single 
most important correctness property for this feature (not attributing a slot to 
the wrong job). What it doesn't cover: (a) the `assignedSlots` empty / 
resource-manager-null degradation from Issue 1; (b) the public 
`build(SeaTunnelServer)` entry point, where `getCoordinatorService()` and the 
two IMap reads actually live; (c) a serializer round-trip for the new class id 
`17`; (d) anything at the REST layer. I verified the precedent cited for (d) is 
real: `Pendi
 ngJobsRestIT.java`, `RealtimeMetricsRestIT.java`, and 
`JobInfoDagStabilityRestIT.java` all exist under 
`seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/`,
 each asserting real endpoint output for a comparable diagnostics endpoint — so 
this codebase does have an established convention of adding a dedicated 
`*RestIT` for a new REST diagnostics surface, this PR just didn't follow it. 
(I'll also note, in fairness, that the older `/running-jobs/summary` endpoint 
itself has no dedicated `RestIT`, so the convention isn't applied 100% 
uniformly — but three separate precedents is enough that I'd still ask for one 
here.)
   - **Potential risk:** a routing regression, a serialization regression on 
class id `17`, or an accidental JSON shape change would all pass CI silently — 
nothing in the current test suite proves the endpoint returns anything over 
HTTP at all.
   - **Best improvement:** add the two missing unit cases (null resource 
manager; `ownedSlotProfiles` populated but `assignedSlots` empty/stale) 
asserting the response is either non-zero or explicitly marked degraded per 
Issue 1's fix; add a `RunningJobSlotUsageRestIT` that submits a job and asserts 
`slotCount` matches its real parallelism, following the `PendingJobsRestIT` 
structure.
   - **Severity:** Low
   
   **Issue 7: `slotCount` counts task-group entries, not distinct physical 
slots — the field/endpoint name implies otherwise**
   - **Location:** `RunningJobSlotUsageBuilder.java:428-433`, `:503-509`
   - **Problem description:** `addSlot` increments once per `TaskGroupLocation` 
entry that passes the filter. If the same physical `SlotProfile` were ever 
referenced by two task groups, or the same pipeline showed up under two 
`PipelineLocation` keys, it would be counted twice, even though `SlotProfile`'s 
own `equals`/`hashCode` (worker + slotID + sequence) would treat those as the 
same slot. The docs word this carefully ("assigned task group slots"), so docs 
and code agree; the mismatch is between the implementation and what a reader 
expects from the endpoint's name and the bare field name `slotCount`.
   - **Potential risk:** capacity-planning conclusions drawn from `slotCount` 
could overstate real slot occupancy in an edge case, and the divergence would 
be invisible because the two numbers agree in the common case.
   - **Best improvement:** either de-duplicate by `SlotKey` before counting (a 
`Set<SlotKey>` per `SlotUsage`, increment only on first insert) so the field 
means what its name says, or keep current semantics and rename to 
`taskGroupSlotCount` so the counting unit is unambiguous. Keep docs in sync 
either way.
   - **Severity:** Low
   
   # 2. Code Quality Assessment
   
   ## 2.1 Coding Standards
   - ASF license headers present on all new source files — verified.
   - No wildcard imports.
   - Every new class carries a class-level Javadoc, and the non-trivial methods 
(`build`, `aggregatePipelineSlots`, `run`, `getRunningJobSlotUsageJson`) have 
explanatory comments — this meets the documentation bar from the project guide. 
Gap: none of the comments state the freshness contract between the two data 
sources being joined; given Issue 1, `build(SeaTunnelServer)` should say in the 
Javadoc that `getAssignedSlots` is a master-local, heartbeat-refreshed snapshot 
that can be empty on a fresh master, and what the response means when that 
happens.
   - Package placement (`server.trace`) is consistent with the existing 
`TaskMappingBuilder` in the same package — no concern there.
   - `GetRunningJobSlotUsageOperation` omits the empty 
`writeInternal`/`readInternal` overrides that `GetRunningJobMetricsOperation` 
carries; functionally irrelevant since both just delegate to `super`, purely 
cosmetic.
   
   ## 2.2 Test Coverage and Test Stability
   - Flaky-pattern check on `RunningJobSlotUsageBuilderTest`: no 
`Thread.sleep`, no `Awaitility`/polling, no network or container dependency, no 
shared static state, no order-sensitive assertions beyond the deterministic 
`TreeMap`/`TreeSet` ordering already used in the production code, no 
floating-point comparisons. Nothing in this test can contribute to CI flakiness.
   - **Test stability rating: Stable**, backed by the above — no High-severity 
flaky-test issue applies here.
   - Coverage gap: as detailed in Issue 6, the untested paths are exactly the 
failure modes, which is the wrong side of the line for a diagnostics API whose 
value proposition is being trustworthy during incidents.
   - Minor, non-blocking observation I hadn't called out before: every 
`SlotProfile` constructed in the test uses the same literal `"test"` for 
`sequence`, so the test's differentiation between "the same slot" and "a 
different slot" relies entirely on `ownerJobId`/`slotId`/`worker`, never on 
`sequence` itself. That's fine for what it's testing, but it means `sequence`'s 
contribution to `SlotKey` uniqueness is exercised only implicitly.
   
   ## 2.3 Documentation Updates
   - All four REST reference files updated: 
`docs/en/engines/zeta/rest-api-v1.md`, `docs/en/engines/zeta/rest-api-v2.md`, 
`docs/zh/engines/zeta/rest-api-v1.md`, `docs/zh/engines/zeta/rest-api-v2.md` — 
re-diffed the English pair myself; both API versions and both languages are 
covered, and the zh content is a real translation, not a placeholder copy.
   - Example payload matches the code: `jobId` is a string 
(`String.valueOf(jobId)`), `slotCount` an int, `pipelineSlotCounts` keyed by 
pipeline id, `workerSlotCounts` keyed by worker address string. The v1 doc 
correctly uses the `/hazelcast/rest/maps/` prefix.
   - The explanatory notes are accurate for the normal path. Missing: the 
meaning of `slotCount: 0` is ambiguous per Issue 1, the counting unit isn't 
spelled out per Issue 7, and the per-request cost characteristic isn't 
mentioned per Issue 4 — these are doc follow-ups to the code fixes, not 
independent doc bugs.
   
   # 3. Architectural Soundness
   
   ## 3.1 Elegance — **Precise fix with one debatable design choice**
   The servlet -> service -> operation -> builder decomposition matches every 
other Zeta REST endpoint's shape, and reusing `BaseServlet`/`BaseService` keeps 
the new code small. Splitting the pure aggregation into a static builder with a 
package-private overload is exactly what makes a mock-free, deterministic unit 
test possible. The one choice I'd push back on is the two-source join itself: 
intersecting a cluster-replicated `IMap` with a master-local in-memory snapshot 
means the answer is only as trustworthy as the weaker source, and the weaker 
source is empty in exactly the failover scenario the endpoint exists to help 
diagnose (Issue 1). A single authoritative source, or an explicit "which 
sources did we manage to read" field, would be a more honest design.
   
   ## 3.2 Maintainability
   Good — small classes, clear names, no inheritance tricks, no shared mutable 
state, easy to unit test in isolation. The main maintainability risk is 
reputational: if operators learn this endpoint sometimes reports unexplained 
zeros, the API loses its diagnostic value. Fixing Issue 1 plus the Issue 5 WARN 
removes that risk.
   
   ## 3.3 Extensibility
   A per-job/per-pipeline/per-worker breakdown is a solid foundation for a 
future Web UI panel or slot-pressure alerting, and the response shape leaves 
room for additional fields (e.g. a degraded flag) without breaking existing 
callers. Once something actually polls it in production, Issues 3 and 4 stop 
being theoretical.
   
   ## 3.4 Historical-Version Compatibility
   No existing REST contract, config option, default value, serialized state, 
checkpoint, or savepoint format is touched; old clients are unaffected. New 
serializer id `17` doesn't collide with current `dev` (max `16`). The only 
rolling-upgrade exposure is the new-node-calls-old-master case noted in 1.2, 
which is confined to this new endpoint and breaks nothing pre-existing. No 
`incompatible-changes.md` entry is required.
   
   # 4. Issue Summary
   
   | # | Issue | Location | Severity |
   | --- | --- | --- | --- |
   | 1 | `slotCount: 0` is indistinguishable from "couldn't read the slot 
source", exactly during failover/cold-start | 
`RunningJobSlotUsageBuilder.java:68-73, 446-448` | Medium |
   | 2 | `getCoordinatorService()` called once per running job in the filter 
lambda; can sleep up to 1.5s per call | `RunningJobSlotUsageBuilder.java:59-66` 
| Medium |
   | 3 | Aggregation runs inline on a shared Hazelcast generic operation 
thread; sibling operation deliberately offloads | 
`GetRunningJobSlotUsageOperation.java:92-96` | Medium |
   | 4 | Full cluster-wide owned-slot map materialized per request, including 
discarded jobs; no cache or rate limit | 
`RunningJobSlotUsageBuilder.java:88-92, 423-426` | Low |
   | 5 | No logging anywhere in the new path; mastership loss on the forwarding 
branch is an unhandled 500 | `RunningJobSlotUsageService.java:35-43` | Low |
   | 6 | No coverage of the degraded path, the public `build(SeaTunnelServer)` 
entry point, serializer round-trip, or REST layer | 
`RunningJobSlotUsageBuilderTest.java` | Low |
   | 7 | `slotCount` counts task-group entries, not distinct slots; name 
implies otherwise | `RunningJobSlotUsageBuilder.java:428-433, 503-509` | Low |
   
   # 5. Merge Recommendation
   
   ### Conclusion: Ready to merge after fixes
   
   **1. Blockers (must fix before merge)**
   - Issue 1 — a diagnostics endpoint must not encode "I couldn't read the 
data" as the same value it uses for "the answer is genuinely zero". Still 
unaddressed since my 2026-08-03/2026-08-16 review of this same code.
   - Issue 2 — hoist `getCoordinatorService()` out of the per-job lambda; 
two-line fix, removes an O(N) x 1.5s blocking amplification.
   - Issue 3 — offload the aggregation off the shared Hazelcast operation 
thread, following the `GetRunningJobMetricsOperation` precedent in the same 
package. Blocking shared operation threads is an engine-wide hazard, not an 
endpoint-local inconvenience.
   
   **2. Recommended fixes (non-blocking)**
   - Issue 4 — narrow the `IMap` read to running-job keys, or add a short-TTL 
cache.
   - Issue 5 — single server lookup, try/catch around the forwarded `.join()`, 
WARN when the slot source is unavailable.
   - Issue 6 — add the two degraded-path unit cases and one REST-level IT 
asserting real slot counts, following `PendingJobsRestIT`.
   - Issue 7 — de-duplicate by `SlotKey`, or rename the field so the counting 
unit is unambiguous.
   
   **Overall assessment**
   
   The feature is worth having and the plumbing is well done: routing is 
correct on both API versions, the serializer id is free, 
`SlotProfile.ownerJobID` really is on the wire so the cross-node intersection 
works, the resource-manager field is properly `volatile`, the endpoint is 
verifiably free of any allocation/deallocation side effect, iteration is 
CME-safe, compatibility is cleanly additive, and all four doc files are updated 
accurately in both languages. Choosing `getInitializedResourceManager()` over 
the lazy-initializing variant was the right instinct for a read-only path.
   
   What's holding it back is that the three blockers cluster around the same 
theme: the endpoint behaves worst exactly when it's needed most — during a 
master transition it can silently return all zeros while blocking a shared 
cluster thread for up to 1.5s per running job, with nothing logged anywhere to 
explain why. Two of the three blockers are a handful of lines each (Issues 2 
and 3); Issue 1 needs one deliberate decision (degraded flag vs. 
coordinator-only fallback) plus a test for it. Once those land, with a 
REST-level test proving real slot counts end to end, this is a solid, useful 
addition to Zeta diagnostics.
   
   **CI status at the time of this review:** `Build` is green (`SUCCESS`) at 
this head — an improvement since 2026-08-16, when the apache-side check was red 
due to an infrastructure-only `paimon-connector-it` Maven-wrapper `429 Too Many 
Requests` failure in the fork run, unrelated to this diff. No CI action needed 
right now beyond the source fixes above. `mergeStateStatus` is `BLOCKED`, which 
on this repo means "awaiting maintainer review", not a merge conflict — 
`mergeable` reports `MERGEABLE`.
   


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