This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12130-ae9132fe0767eee612259c293dc28b73a48dfae4 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
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); + } +}
