This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 41a9c4b8fb [Improve][Zeta] Apply the finished-jobs page before the
per-job metrics and DAG lookups (#12130)
41a9c4b8fb is described below
commit 41a9c4b8fb0bbd6e0538cfd2728ead4651032ff6
Author: Sudarshan Kumar Kaushik <[email protected]>
AuthorDate: Mon Sep 14 15:30:05 2026 +0000
[Improve][Zeta] Apply the finished-jobs page before the per-job metrics and
DAG lookups (#12130)
---
docs/en/engines/zeta/rest-api-v2.md | 12 +-
.../introduction/concepts/incompatible-changes.md | 16 ++
docs/zh/engines/zeta/rest-api-v2.md | 10 +-
.../introduction/concepts/incompatible-changes.md | 7 +
.../engine/server/rest/service/JobInfoService.java | 85 ++++++++--
.../server/rest/servlet/FinishedJobsServlet.java | 21 ++-
.../server/rest/servlet/PageBaseServlet.java | 145 +++++++++++++----
.../rest/service/JobInfoServiceNullSafetyTest.java | 48 ++++++
.../rest/servlet/FinishedJobsServletTest.java | 179 +++++++++++++++++++++
9 files changed, 477 insertions(+), 46 deletions(-)
diff --git a/docs/en/engines/zeta/rest-api-v2.md
b/docs/en/engines/zeta/rest-api-v2.md
index efe4d04a3e..8616c953fd 100644
--- a/docs/en/engines/zeta/rest-api-v2.md
+++ b/docs/en/engines/zeta/rest-api-v2.md
@@ -783,8 +783,16 @@ When we can't get the job info, the response will be:
> | name | type | data type | description
> |
> |-------|----------|-----------|-----------------------------------------------------------------------------------|
> | state | optional | string | finished job status.
> `FINISHED`,`CANCELED`,`FAILED`,`SAVEPOINT_DONE`,`UNKNOWABLE` |
-> | page | optional | int | page number.
|
-> | rows | optional | int | page size.
|
+> | page | optional | int | page number. Must be an integer greater
than 0. |
+> | rows | optional | int | page size, defaults to 10. Must be an
integer greater than 0. |
+
+When `page` is supplied, the response is wrapped as `{"data": [...], "total":
n}`, where `total`
+is the number of jobs matching `state` before the page is applied. When it is
omitted, the bare
+array is returned.
+
+A `page` or `rows` value that is not an integer, or is not greater than 0,
returns `400`. A page
+starting beyond the end of the result set also returns `400`, while a page
starting exactly at
+`total` returns an empty page.
#### Responses
diff --git a/docs/en/introduction/concepts/incompatible-changes.md
b/docs/en/introduction/concepts/incompatible-changes.md
index 42ae00db16..0dc9cb8392 100644
--- a/docs/en/introduction/concepts/incompatible-changes.md
+++ b/docs/en/introduction/concepts/incompatible-changes.md
@@ -5,6 +5,22 @@ You need to check this document before you upgrade to related
version.
## dev
+### Zeta REST Pagination Parameter Validation
+
+- **Behavior change: `page` and `rows` are validated on paginated endpoints**
+ - **Affected component**: `seatunnel-engine-server`, REST endpoints `GET
/finished-jobs/:state`,
+ `GET /running-jobs` and `GET /running-jobs/summary`. The latter two are
served by the same
+ `RunningJobsServlet` instance, so both receive the validation.
+ - **Description**: These endpoints now reject a `page` or `rows` value that
is not an integer or
+ is not greater than 0, and reject a page whose start offset would overflow
a 32-bit integer.
+ Previously `rows=0` was accepted and returned an empty page, a negative
`rows` produced an
+ internal error, and a sufficiently large `page` combined with `rows` could
wrap to a small
+ positive offset and silently return the wrong page.
+ - **Impact**: Requests that relied on `rows=0` returning an empty page now
receive `400` with a
+ message naming the offending parameter. Callers passing valid positive
values are unaffected.
+ The response shape, the `{"data": [...], "total": n}` envelope, and the
behaviour of a page
+ starting exactly at `total`, which still returns an empty page, are all
unchanged.
+
### MySQL CDC Schema-Change Parsing
- **Behavior change: DDL parser listener errors are propagated**
diff --git a/docs/zh/engines/zeta/rest-api-v2.md
b/docs/zh/engines/zeta/rest-api-v2.md
index c2e4fabd88..9ef191f636 100644
--- a/docs/zh/engines/zeta/rest-api-v2.md
+++ b/docs/zh/engines/zeta/rest-api-v2.md
@@ -756,8 +756,14 @@ seatunnel:
> | 参数名称 | 是否必传 | 参数类型 | 参数描述
> |
> |-------|----------|--------|-----------------------------------------------------------------------------------|
> | state | optional | string | finished job status.
> `FINISHED`,`CANCELED`,`FAILED`,`SAVEPOINT_DONE`,`UNKNOWABLE` |
-> | page | 否 | int | 页号 |
-> | rows | 否 | int | 每页行数 |
+> | page | 否 | int | 页号,必须是大于 0 的整数 |
+> | rows | 否 | int | 每页行数,默认为 10,必须是大于 0 的整数 |
+
+当传入 `page` 时,响应会被包装为 `{"data": [...], "total": n}`,其中 `total` 是分页之前匹配
+`state` 的作业总数。未传入 `page` 时,直接返回数组。
+
+`page` 或 `rows` 不是整数,或者不大于 0 时,返回 `400`。起始位置超出结果集末尾时同样返回 `400`,
+而起始位置恰好等于 `total` 时返回空页。
#### 响应
diff --git a/docs/zh/introduction/concepts/incompatible-changes.md
b/docs/zh/introduction/concepts/incompatible-changes.md
index d653c20908..21a5217a28 100644
--- a/docs/zh/introduction/concepts/incompatible-changes.md
+++ b/docs/zh/introduction/concepts/incompatible-changes.md
@@ -4,6 +4,13 @@
## dev
+### Zeta REST 分页参数校验
+
+- **行为变更:分页接口开始校验 `page` 与 `rows`**
+ - **影响范围**:`seatunnel-engine-server`,REST 接口 `GET
/finished-jobs/:state`、`GET /running-jobs` 与 `GET
/running-jobs/summary`。后两者由同一个 `RunningJobsServlet` 实例提供服务,因此都会受到该校验。
+ - **变更说明**:这些接口现在会拒绝非整数或不大于 0 的 `page` 与 `rows`,并拒绝起始偏移量会超出 32 位整数范围的分页请求。此前
`rows=0` 会被接受并返回空页,负数 `rows` 会引发内部错误,而足够大的 `page` 与 `rows`
组合可能溢出为一个较小的正偏移量,从而静默返回错误的页。
+ - **影响**:依赖 `rows=0` 返回空页的请求现在会收到
`400`,错误信息中会指明具体参数。传入合法正整数的调用方不受影响。响应结构、`{"data": [...], "total": n}`
包装格式,以及起始位置恰好等于 `total` 时仍返回空页的行为,均保持不变。
+
### MySQL CDC Schema-Change 解析
- **行为变更:向上传播 DDL 解析监听器错误**
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/JobInfoService.java
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/JobInfoService.java
index bb73f58c47..5a9707e2f1 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/JobInfoService.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/JobInfoService.java
@@ -58,6 +58,7 @@ import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
import static
org.apache.seatunnel.engine.server.rest.RestConstant.CONFIG_FORMAT;
@@ -115,12 +116,55 @@ public class JobInfoService extends BaseService {
}
public JsonArray getJobsByStateJson(String state) {
- IMap<Long, JobState> finishedJob =
-
nodeEngine.getHazelcastInstance().getMap(Constant.IMAP_FINISHED_JOB_STATE);
+ IMap<Long, JobDAGInfo> finishedJobDAGInfo =
+
nodeEngine.getHazelcastInstance().getMap(Constant.IMAP_FINISHED_JOB_VERTEX_INFO);
+
+ return matchingJobStates(state).stream()
+ .map(jobState -> toJobInfoJson(jobState, finishedJobDAGInfo))
+ .collect(JsonArray::new, JsonArray::add, JsonArray::add);
+ }
+
+ /**
+ * Returns one page of finished jobs in the given state, newest first,
along with the total
+ * number of matches before slicing.
+ *
+ * <p>The slice is applied before the per-job metrics and DAG lookups, so
a paged request pays
+ * those two point lookups and the metric aggregation only for the rows it
actually returns
+ * rather than for every retained job.
+ */
+ public JobPage getJobsByStateJson(String state, int start, int rows) {
+ // checked here rather than left to Stream.skip/limit, so any caller
gets a message naming
+ // the argument instead of an exception from inside the stream pipeline
+ if (state == null) {
+ throw new IllegalArgumentException("state must not be null, use
\"\" for any state");
+ }
+ if (start < 0) {
+ throw new IllegalArgumentException("start must not be negative,
but was: " + start);
+ }
+ if (rows < 1) {
+ throw new IllegalArgumentException("rows must be greater than 0,
but was: " + rows);
+ }
IMap<Long, JobDAGInfo> finishedJobDAGInfo =
nodeEngine.getHazelcastInstance().getMap(Constant.IMAP_FINISHED_JOB_VERTEX_INFO);
+ List<JobState> matches = matchingJobStates(state);
+
+ JsonArray data =
+ matches.stream()
+ .skip(start)
+ .limit(rows)
+ .map(jobState -> toJobInfoJson(jobState,
finishedJobDAGInfo))
+ .collect(JsonArray::new, JsonArray::add,
JsonArray::add);
+
+ return new JobPage(data, matches.size());
+ }
+
+ /** Finished jobs matching {@code state} (all of them when it is empty),
newest first. */
+ private List<JobState> matchingJobStates(String state) {
+ IMap<Long, JobState> finishedJob =
+
nodeEngine.getHazelcastInstance().getMap(Constant.IMAP_FINISHED_JOB_STATE);
+
return finishedJob.values().stream()
.filter(java.util.Objects::nonNull)
.filter(
@@ -134,15 +178,34 @@ public class JobInfoService extends BaseService {
Comparator.comparing(
JobState::getFinishTime,
Comparator.nullsLast(Comparator.reverseOrder())))
- .map(
- jobState -> {
- Long jobId = jobState.getJobId();
- return getJobInfoJson(
- jobState,
- getFinishedJobMetricsJson(jobId),
- getFinishedJobDAGInfo(finishedJobDAGInfo,
jobId));
- })
- .collect(JsonArray::new, JsonArray::add, JsonArray::add);
+ .collect(Collectors.toList());
+ }
+
+ private JsonObject toJobInfoJson(JobState jobState, IMap<Long, JobDAGInfo>
finishedJobDAGInfo) {
+ Long jobId = jobState.getJobId();
+ return getJobInfoJson(
+ jobState,
+ getFinishedJobMetricsJson(jobId),
+ getFinishedJobDAGInfo(finishedJobDAGInfo, jobId));
+ }
+
+ /** A page of job listings plus the number of matches before the page was
taken. */
+ public static final class JobPage {
+ private final JsonArray data;
+ private final int total;
+
+ public JobPage(JsonArray data, int total) {
+ this.data = data;
+ this.total = total;
+ }
+
+ public JsonArray getData() {
+ return data;
+ }
+
+ public int getTotal() {
+ return total;
+ }
}
private String getFinishedJobMetricsJson(Long jobId) {
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/FinishedJobsServlet.java
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/FinishedJobsServlet.java
index dc1a8aefa6..7c53f68db3 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/FinishedJobsServlet.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/FinishedJobsServlet.java
@@ -34,8 +34,13 @@ public class FinishedJobsServlet extends PageBaseServlet {
private final JobInfoService jobInfoService;
public FinishedJobsServlet(NodeEngineImpl nodeEngine) {
+ this(nodeEngine, new JobInfoService(nodeEngine));
+ }
+
+ /** Visible for testing, so the service can be substituted without a
running node engine. */
+ FinishedJobsServlet(NodeEngineImpl nodeEngine, JobInfoService
jobInfoService) {
super(nodeEngine);
- this.jobInfoService = new JobInfoService(nodeEngine);
+ this.jobInfoService = jobInfoService;
}
@Override
@@ -50,6 +55,18 @@ public class FinishedJobsServlet extends PageBaseServlet {
state = "";
}
- writeJsonWithPagination(req, resp,
jobInfoService.getJobsByStateJson(state));
+ PageParams pageParams = pageParams(req);
+ if (pageParams == null) {
+ writeJson(resp, jobInfoService.getJobsByStateJson(state));
+ return;
+ }
+
+ // slice at the source: building the whole listing and paginating
afterwards would do the
+ // per-job metrics and DAG lookups for every retained job just to
discard all but one page
+ JobInfoService.JobPage jobPage =
+ jobInfoService.getJobsByStateJson(
+ state, pageParams.getStart(), pageParams.getRows());
+ checkPageInRange(pageParams, jobPage.getTotal());
+ writeJsonPage(resp, jobPage.getData(), jobPage.getTotal());
}
}
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/PageBaseServlet.java
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/PageBaseServlet.java
index 147dfe0100..b4260c3ec5 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/PageBaseServlet.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/servlet/PageBaseServlet.java
@@ -28,6 +28,8 @@ import java.io.IOException;
import java.util.Map;
public class PageBaseServlet extends BaseServlet {
+ private static final int DEFAULT_ROWS = 10;
+
private final String pageParam = "page";
private final String rowsParam = "rows";
@@ -35,41 +37,126 @@ public class PageBaseServlet extends BaseServlet {
super(nodeEngine);
}
+ /**
+ * Paginates an already-built array, writing {@code {"data": [], "total":
n}} when a page was
+ * requested and the bare array otherwise.
+ *
+ * <p>Endpoints that can slice at the source should do so and use {@link
#writeJsonPage}
+ * instead, because this method has already paid the cost of building
every row before
+ * discarding all but one page.
+ */
protected void writeJsonWithPagination(
HttpServletRequest req, HttpServletResponse resp, JsonArray
jsonArray)
throws IOException {
int total = jsonArray.size();
- // fetch pagination params, if page exist, then paginate
data,pagination data format like:
- // {"data": [], "total": 10}
- Map<String, String> parameterMap = getParameterMap(req);
- if (parameterMap != null && parameterMap.containsKey(pageParam)) {
- int page = Integer.parseInt(parameterMap.get(pageParam));
- int rows =
- parameterMap.get(rowsParam) != null
- ? Integer.parseInt(parameterMap.get(rowsParam))
- : 10;
- int start = (page - 1) * rows;
- if (start > total || page < 1) {
- throw new IllegalArgumentException(
- page < 1
- ? "Page number must be greater than 0"
- : "Page number exceeds total pages");
- }
- JsonArray paginatedArray = new JsonArray();
- jsonArray
- .values()
- .subList(start, Math.min(start + rows, total))
- .forEach(
- t -> {
- paginatedArray.add(t);
- });
- JsonObject paginatedObj = new JsonObject();
- paginatedObj.add("data", paginatedArray);
- paginatedObj.add("total", total);
- writeJson(resp, paginatedObj);
- } else {
+ PageParams pageParams = pageParams(req);
+ if (pageParams == null) {
writeJson(resp, jsonArray);
+ return;
+ }
+ checkPageInRange(pageParams, total);
+
+ int start = pageParams.getStart();
+ JsonArray paginatedArray = new JsonArray();
+ jsonArray
+ .values()
+ .subList(start, Math.min(start + pageParams.getRows(), total))
+ .forEach(
+ t -> {
+ paginatedArray.add(t);
+ });
+ writeJsonPage(resp, paginatedArray, total);
+ }
+
+ /**
+ * A requested page. Absent when the caller sent no {@code page} parameter.
+ *
+ * <p>The start offset is computed and validated in {@link
#pageParams(HttpServletRequest)}
+ * rather than derived on demand, so that no accessor here can throw.
+ */
+ protected static final class PageParams {
+ private final int rows;
+ private final int start;
+
+ private PageParams(int rows, int start) {
+ this.rows = rows;
+ this.start = start;
}
+
+ public int getRows() {
+ return rows;
+ }
+
+ public int getStart() {
+ return start;
+ }
+ }
+
+ /**
+ * Reads and validates the pagination parameters, or returns {@code null}
when the request does
+ * not ask for a page.
+ *
+ * <p>This is the single definition of pagination input handling. {@link
+ * #writeJsonWithPagination} routes through it as well, so every paginated
endpoint accepts and
+ * rejects the same input.
+ *
+ * @throws IllegalArgumentException if the values are unparseable,
non-positive, or describe an
+ * offset too large to address, all of which the servlet layer reports
as a 400
+ */
+ protected PageParams pageParams(HttpServletRequest req) {
+ Map<String, String> parameterMap = getParameterMap(req);
+ if (parameterMap == null || !parameterMap.containsKey(pageParam)) {
+ return null;
+ }
+ int page = positiveIntParam(parameterMap.get(pageParam), pageParam);
+ String rowsValue = parameterMap.get(rowsParam);
+ int rows = rowsValue != null ? positiveIntParam(rowsValue, rowsParam)
: DEFAULT_ROWS;
+
+ // widened before multiplying: page and rows are both caller
controlled, and an int
+ // overflow here can wrap to a small positive offset that passes
checkPageInRange and
+ // silently serves the wrong page
+ long start = (long) (page - 1) * rows;
+ if (start > Integer.MAX_VALUE) {
+ throw new IllegalArgumentException("Page number exceeds total
pages");
+ }
+ return new PageParams(rows, (int) start);
+ }
+
+ private int positiveIntParam(String value, String name) {
+ int parsed;
+ try {
+ parsed = Integer.parseInt(value.trim());
+ } catch (NumberFormatException e) {
+ // the JDK message quotes the raw input without naming the
parameter
+ throw new IllegalArgumentException(
+ "Parameter '" + name + "' must be an integer, but was: " +
value);
+ }
+ if (parsed < 1) {
+ throw new IllegalArgumentException("Parameter '" + name + "' must
be greater than 0");
+ }
+ return parsed;
+ }
+
+ /**
+ * Rejects a page that starts past the end of the result set.
+ *
+ * <p>A page starting exactly at {@code total} is allowed and yields an
empty page. That is
+ * deliberate: it is the behaviour the pre-existing pagination has always
had, and changing it
+ * would alter the response for callers that walk to the end of a listing.
+ */
+ protected void checkPageInRange(PageParams pageParams, int total) {
+ if (pageParams.getStart() > total) {
+ throw new IllegalArgumentException("Page number exceeds total
pages");
+ }
+ }
+
+ /** Writes the {@code {"data": [], "total": n}} envelope for an
already-sliced page. */
+ protected void writeJsonPage(HttpServletResponse resp, JsonArray pageData,
int total)
+ throws IOException {
+ JsonObject paginatedObj = new JsonObject();
+ paginatedObj.add("data", pageData);
+ paginatedObj.add("total", total);
+ writeJson(resp, paginatedObj);
}
}
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/service/JobInfoServiceNullSafetyTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/service/JobInfoServiceNullSafetyTest.java
index 2fd2fb6307..09b89097b2 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/service/JobInfoServiceNullSafetyTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/service/JobInfoServiceNullSafetyTest.java
@@ -42,10 +42,15 @@ import com.hazelcast.spi.impl.NodeEngineImpl;
import java.io.IOException;
import java.lang.reflect.Field;
import java.lang.reflect.Method;
+import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
+import java.util.List;
+import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
class JobInfoServiceNullSafetyTest {
@@ -104,6 +109,49 @@ class JobInfoServiceNullSafetyTest {
Assertions.assertEquals(createTime, getFieldValue(basicInfo,
"createTime"));
}
+ /**
+ * A paged request must pay the per-job metrics and DAG lookups for the
rows it returns, not for
+ * every retained job. Verifying the lookup count is the point of the
test: asserting only on
+ * the returned JSON would still pass if the whole listing were built and
then sliced.
+ */
+ @Test
+ void shouldApplyPageBeforePerJobLookups() {
+ List<Object> storedStates = new ArrayList<>();
+ for (int index = 0; index < 25; index++) {
+ storedStates.add(buildJobState((long) index, 1000L, 2000L +
index));
+ }
+ when(finishedJobStateMap.values()).thenReturn(storedStates);
+ when(finishedJobMetricsMap.getOrDefault(any(),
any())).thenReturn(JobMetrics.empty());
+ when(finishedJobVertexInfoMap.get(any())).thenReturn(null);
+
+ JobInfoService.JobPage page = jobInfoService.getJobsByStateJson("", 0,
10);
+
+ Assertions.assertEquals(10, page.getData().size());
+ Assertions.assertEquals(25, page.getTotal(), "total must count matches
before slicing");
+ verify(finishedJobMetricsMap, times(10)).getOrDefault(any(), any());
+ verify(finishedJobVertexInfoMap, times(10)).get(any());
+ }
+
+ /** Newest finish time first, so that paging over a growing history stays
stable. */
+ @Test
+ void shouldReturnNewestFinishedJobsFirst() {
+ List<Object> storedStates = new ArrayList<>();
+ for (int index = 0; index < 5; index++) {
+ storedStates.add(buildJobState((long) index, 1000L, 2000L +
index));
+ }
+ when(finishedJobStateMap.values()).thenReturn(storedStates);
+ when(finishedJobMetricsMap.getOrDefault(any(),
any())).thenReturn(JobMetrics.empty());
+ when(finishedJobVertexInfoMap.get(any())).thenReturn(null);
+
+ JobInfoService.JobPage page = jobInfoService.getJobsByStateJson("", 0,
2);
+
+ Assertions.assertEquals(2, page.getData().size());
+ Assertions.assertEquals(
+ "4",
page.getData().get(0).asObject().getString(RestConstant.JOB_ID, null));
+ Assertions.assertEquals(
+ "3",
page.getData().get(1).asObject().getString(RestConstant.JOB_ID, null));
+ }
+
private JobHistoryService.JobState buildJobState(Long jobId, Long
startTime, Long finishTime) {
return new JobHistoryService.JobState(
jobId,
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/servlet/FinishedJobsServletTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/servlet/FinishedJobsServletTest.java
new file mode 100644
index 0000000000..3a0b369a1c
--- /dev/null
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/rest/servlet/FinishedJobsServletTest.java
@@ -0,0 +1,179 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.engine.server.rest.servlet;
+
+import org.apache.seatunnel.engine.server.rest.service.JobInfoService;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import com.hazelcast.internal.json.Json;
+import com.hazelcast.internal.json.JsonArray;
+import com.hazelcast.internal.json.JsonObject;
+import com.hazelcast.spi.impl.NodeEngineImpl;
+
+import javax.servlet.http.HttpServletRequest;
+import javax.servlet.http.HttpServletResponse;
+
+import java.io.PrintWriter;
+import java.io.StringWriter;
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Covers the servlet side of finished-job pagination: parameter validation,
the out-of-range
+ * boundary, and the routing between the paged and unpaged service calls.
+ */
+class FinishedJobsServletTest {
+
+ private JobInfoService jobInfoService;
+ private FinishedJobsServlet servlet;
+ private HttpServletRequest request;
+ private HttpServletResponse response;
+ private StringWriter output;
+
+ @BeforeEach
+ void setUp() throws Exception {
+ NodeEngineImpl nodeEngine = mock(NodeEngineImpl.class);
+ jobInfoService = mock(JobInfoService.class);
+ servlet = new FinishedJobsServlet(nodeEngine, jobInfoService);
+
+ request = mock(HttpServletRequest.class);
+ response = mock(HttpServletResponse.class);
+ output = new StringWriter();
+ when(response.getWriter()).thenReturn(new PrintWriter(output));
+ when(request.getPathInfo()).thenReturn(null);
+ }
+
+ private void withParams(String... keyValuePairs) {
+ Map<String, String[]> params = new HashMap<>();
+ for (int index = 0; index < keyValuePairs.length; index += 2) {
+ params.put(keyValuePairs[index], new String[] {keyValuePairs[index
+ 1]});
+ }
+ when(request.getParameterMap()).thenReturn(params);
+ }
+
+ /** Without a page parameter the caller still gets the whole listing as a
bare array. */
+ @Test
+ void shouldReturnBareArrayWhenNoPageRequested() throws Exception {
+ withParams();
+ when(jobInfoService.getJobsByStateJson("")).thenReturn(new
JsonArray().add("job"));
+
+ servlet.doGet(request, response);
+
+ Assertions.assertEquals("[\"job\"]", output.toString());
+ verify(jobInfoService, never()).getJobsByStateJson(anyString(),
anyInt(), anyInt());
+ }
+
+ /** With a page parameter the servlet must use the slicing overload and
the envelope. */
+ @Test
+ void shouldUsePagedOverloadAndWriteEnvelopeWhenPageRequested() throws
Exception {
+ withParams("page", "1", "rows", "2");
+ when(jobInfoService.getJobsByStateJson("", 0, 2))
+ .thenReturn(page(new JsonArray().add("a").add("b"), 7));
+
+ servlet.doGet(request, response);
+
+ JsonObject written = Json.parse(output.toString()).asObject();
+ Assertions.assertEquals(2, written.get("data").asArray().size());
+ Assertions.assertEquals(7, written.get("total").asInt());
+ verify(jobInfoService, never()).getJobsByStateJson(anyString());
+ }
+
+ /** A page starting exactly at total is allowed and yields an empty page,
matching legacy. */
+ @Test
+ void shouldReturnEmptyPageWhenStartEqualsTotal() throws Exception {
+ withParams("page", "3", "rows", "5");
+ when(jobInfoService.getJobsByStateJson("", 10, 5)).thenReturn(page(new
JsonArray(), 10));
+
+ servlet.doGet(request, response);
+
+ JsonObject written = Json.parse(output.toString()).asObject();
+ Assertions.assertEquals(0, written.get("data").asArray().size());
+ Assertions.assertEquals(10, written.get("total").asInt());
+ }
+
+ @Test
+ void shouldRejectPageStartingBeyondTotal() throws Exception {
+ withParams("page", "4", "rows", "5");
+ when(jobInfoService.getJobsByStateJson("", 15, 5)).thenReturn(page(new
JsonArray(), 10));
+
+ Assertions.assertEquals("Page number exceeds total pages",
assertRejected().getMessage());
+ }
+
+ @Test
+ void shouldRejectZeroRows() {
+ withParams("page", "1", "rows", "0");
+
+ Assertions.assertEquals(
+ "Parameter 'rows' must be greater than 0",
assertRejected().getMessage());
+ }
+
+ @Test
+ void shouldRejectNegativeRows() {
+ withParams("page", "1", "rows", "-5");
+
+ Assertions.assertEquals(
+ "Parameter 'rows' must be greater than 0",
assertRejected().getMessage());
+ }
+
+ @Test
+ void shouldRejectNonPositivePage() {
+ withParams("page", "0");
+
+ Assertions.assertEquals(
+ "Parameter 'page' must be greater than 0",
assertRejected().getMessage());
+ }
+
+ @Test
+ void shouldRejectNonNumericInputWithAMessageNamingTheParameter() {
+ withParams("page", "abc");
+
+ Assertions.assertEquals(
+ "Parameter 'page' must be an integer, but was: abc",
assertRejected().getMessage());
+ }
+
+ /**
+ * The offset is computed in long arithmetic, so a page and row count
whose product overflows an
+ * int is rejected rather than wrapping to a small positive offset and
serving the wrong page.
+ */
+ @Test
+ void shouldRejectOffsetThatOverflowsAnInt() throws Exception {
+ withParams("page", String.valueOf(Integer.MAX_VALUE), "rows", "10");
+
+ Assertions.assertEquals("Page number exceeds total pages",
assertRejected().getMessage());
+ verify(jobInfoService, never()).getJobsByStateJson(anyString(),
anyInt(), anyInt());
+ }
+
+ private IllegalArgumentException assertRejected() {
+ return Assertions.assertThrows(
+ IllegalArgumentException.class, () -> servlet.doGet(request,
response));
+ }
+
+ private JobInfoService.JobPage page(JsonArray data, int total) {
+ return new JobInfoService.JobPage(data, total);
+ }
+}