This is an automated email from the ASF dual-hosted git repository.
davidzollo 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 47223dcd1c [Fix][Zeta] Resolve each member's bound REST HTTP port
(#12298)
47223dcd1c is described below
commit 47223dcd1c96763bf42590073f1c767ce628f402
Author: SEZ <[email protected]>
AuthorDate: Sun Sep 27 19:10:44 2026 +0800
[Fix][Zeta] Resolve each member's bound REST HTTP port (#12298)
---
docs/en/engines/zeta/rest-api-v2.md | 6 +-
docs/zh/engines/zeta/rest-api-v2.md | 6 +-
.../org/apache/seatunnel/engine/e2e/RestApiIT.java | 97 +++++++++++++++++++---
.../seatunnel/engine/server/JettyService.java | 39 ++++++---
.../seatunnel/engine/server/SeaTunnelServer.java | 10 ++-
.../server/operation/GetNodeHttpPortOperation.java | 6 +-
.../server/rest/service/LoggerLevelService.java | 6 +-
.../operation/GetNodeHttpPortOperationTest.java | 68 ++++++++++++++-
8 files changed, 200 insertions(+), 38 deletions(-)
diff --git a/docs/en/engines/zeta/rest-api-v2.md
b/docs/en/engines/zeta/rest-api-v2.md
index f49243e325..26382006fa 100644
--- a/docs/en/engines/zeta/rest-api-v2.md
+++ b/docs/en/engines/zeta/rest-api-v2.md
@@ -78,7 +78,7 @@ seatunnel:
## Web UI and Port 8080 Troubleshooting
- If `http://<host>:8080/` is unreachable, first check whether
`seatunnel.engine.http.enable-http` or `enable-https` is actually enabled. The
`network.rest-api.enabled` setting in `hazelcast.yaml` does not replace the
Jetty switch.
-- If `enable-dynamic-port = true`, the actual listening port may not be 8080.
Jetty will choose the first available port between `port` and `port +
port-range`. Use the startup log `SeaTunnel REST service will start on port
xxx` as the source of truth.
+- If HTTP and `enable-dynamic-port = true` are enabled, the actual listening
port may not be 8080. Jetty chooses the first available port between `port` and
`port + port-range`. Use the Jetty startup log `SeaTunnel REST service started
on http port xxx` as the source of truth. `/logs` and `/loggers?scope=cluster`
resolve and report each member's actual bound HTTP port. The configured `port`
remains unchanged, including when members share an HTTP configuration object.
- If `context-path = /seatunnel`, both the Web UI and REST endpoints move
under that prefix. For example, the overview endpoint becomes
`/seatunnel/overview`.
- The Web UI static resources and REST endpoints share the same Jetty service.
If Jetty does not start, both are unavailable together.
@@ -1543,8 +1543,8 @@ With `?scope=cluster` the answer is one entry per member:
`status` is `SUCCESS` when every member answered, `PARTIAL_FAILURE` when some
did not, and `FAILURE`
when none did; the member that failed carries its own `status` and `error`. A
cluster request reaches
-every member on the REST port of its configuration, so it does not reach
members that took a
-different port through `enable-dynamic-port`.
+every member on its actual bound REST HTTP port, including members that
selected a different
+port through `enable-dynamic-port`. Each member must have HTTP enabled and be
reachable.
</details>
diff --git a/docs/zh/engines/zeta/rest-api-v2.md
b/docs/zh/engines/zeta/rest-api-v2.md
index 1e88085dce..39d39928b9 100644
--- a/docs/zh/engines/zeta/rest-api-v2.md
+++ b/docs/zh/engines/zeta/rest-api-v2.md
@@ -73,7 +73,7 @@ seatunnel:
## Web UI 与 8080 排查
- 如果 `http://<host>:8080/` 打不开,先检查 `seatunnel.engine.http.enable-http` 或
`enable-https` 是否真的开启;仅配置 `hazelcast.yaml` 中的 `network.rest-api.enabled` 不能替代
Jetty 开关。
-- 如果开启了 `enable-dynamic-port = true`,实际监听端口可能不是 8080,而是 `port` 到 `port +
port-range` 之间的第一个空闲端口。以启动日志 `SeaTunnel REST service will start on port xxx` 为准。
+- 如果同时开启 HTTP 和 `enable-dynamic-port = true`,实际监听端口可能不是 8080,而是 `port` 到 `port
+ port-range` 之间的第一个空闲端口。以 Jetty 启动日志 `SeaTunnel REST service started on http
port xxx` 为准。`/logs` 和 `/loggers?scope=cluster` 会解析并报告各节点实际绑定的 HTTP 端口。配置中的
`port` 保持不变,即使多个节点共享同一个 HTTP 配置对象也不例外。
- 如果配置了 `context-path = /seatunnel`,Web UI 首页和 REST 路径都会整体前移,例如概览接口会变成
`/seatunnel/overview`。
- Web UI 静态资源和 REST API 共用同一个 Jetty 服务。只要 Jetty 没启动,两者都会一起不可用。
@@ -1519,8 +1519,8 @@ curl --location
'http://127.0.0.1:8080/submit-job/upload?restoreMode=CHECKPOINT&
```
所有节点都返回结果时 `status` 为 `SUCCESS`,部分节点失败时为 `PARTIAL_FAILURE`,全部失败时为
-`FAILURE`;失败的节点会带上自己的 `status` 与 `error`。集群请求按各节点配置中的 REST 端口访问,因此无法
-访问通过 `enable-dynamic-port` 使用了其它端口的节点。
+`FAILURE`;失败的节点会带上自己的 `status` 与 `error`。集群请求按各节点实际绑定的 REST HTTP 端口访问,
+包括通过 `enable-dynamic-port` 选择了其它端口的节点。各节点都需要启用 HTTP,且其端口可访问。
</details>
diff --git
a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java
b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java
index a25362a965..4e67a4f633 100644
---
a/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java
+++
b/seatunnel-e2e/seatunnel-engine-e2e/connector-seatunnel-e2e-base/src/test/java/org/apache/seatunnel/engine/e2e/RestApiIT.java
@@ -31,6 +31,7 @@ import
org.apache.seatunnel.engine.server.SeaTunnelServerStarter;
import org.apache.seatunnel.engine.server.checkpoint.CheckpointCloseReason;
import
org.apache.seatunnel.engine.server.checkpoint.monitor.CheckpointMonitorService;
import org.apache.seatunnel.engine.server.rest.RestConstant;
+import org.apache.seatunnel.engine.server.rest.service.LogService;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.core.LoggerContext;
@@ -50,21 +51,26 @@ import io.restassured.common.mapper.TypeRef;
import lombok.extern.slf4j.Slf4j;
import java.lang.reflect.Field;
+import java.nio.file.Path;
import java.nio.file.Paths;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
+import java.util.Set;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicReference;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
+import java.util.stream.Collectors;
import static io.restassured.RestAssured.given;
import static
org.apache.seatunnel.e2e.common.util.ContainerUtil.PROJECT_ROOT_PATH;
import static
org.apache.seatunnel.engine.server.rest.RestConstant.CONTEXT_PATH;
+import static org.hamcrest.Matchers.containsInAnyOrder;
import static org.hamcrest.Matchers.containsString;
import static org.hamcrest.Matchers.equalTo;
import static org.hamcrest.Matchers.hasItem;
@@ -108,6 +114,7 @@ public class RestApiIT {
node1Config = ConfigProvider.locateAndGetSeaTunnelConfig();
node1Config.getEngineConfig().getHttpConfig().setPort(8080);
node1Config.getEngineConfig().getHttpConfig().setEnabled(true);
+
node1Config.getEngineConfig().getHttpConfig().setEnableDynamicPort(true);
node1Config.getHazelcastConfig().setClusterName(testClusterName);
node1Config.getEngineConfig().getSlotServiceConfig().setDynamicSlot(false);
node1Config.getEngineConfig().getSlotServiceConfig().setSlotNum(20);
@@ -120,9 +127,9 @@ public class RestApiIT {
node2Tags.setAttribute("node", "node2");
Config node2hzconfig =
node1Config.getHazelcastConfig().setMemberAttributeConfig(node2Tags);
node2Config = ConfigProvider.locateAndGetSeaTunnelConfig();
- // Dynamically generated port
-
node2Config.getEngineConfig().getHttpConfig().setEnableDynamicPort(true);
- node2Config.getEngineConfig().getHttpConfig().setEnabled(true);
+ // Both members deliberately share the same configured port and
mutable HTTP bean.
+ // Node2 must bind another port without changing what either member
advertises.
+
node2Config.getEngineConfig().setHttpConfig(node1Config.getEngineConfig().getHttpConfig());
node2Config.getEngineConfig().getSlotServiceConfig().setDynamicSlot(false);
node2Config.getEngineConfig().getSlotServiceConfig().setSlotNum(20);
node2Config.setHazelcastConfig(node2hzconfig);
@@ -162,12 +169,8 @@ public class RestApiIT {
Assertions.assertEquals(
JobStatus.FINISHED,
batchJobProxy.getJobStatus()));
ports = new HashMap<>();
- ports.put(
- node1.getCluster().getLocalMember().getAddress().getPort(),
- node1Config.getEngineConfig().getHttpConfig().getPort());
- ports.put(
- node2.getCluster().getLocalMember().getAddress().getPort(),
- node2Config.getEngineConfig().getHttpConfig().getPort());
+ ports.put(node1.getCluster().getLocalMember().getAddress().getPort(),
httpPort(node1));
+ ports.put(node2.getCluster().getLocalMember().getAddress().getPort(),
httpPort(node2));
}
@Test
@@ -247,11 +250,61 @@ public class RestApiIT {
}));
}
+ @Test
+ public void testDynamicHttpPortIsResolvableByPeers() throws Exception {
+ int node1HttpPort = httpPort(node1);
+ int node2HttpPort = httpPort(node2);
+
+ Assertions.assertSame(
+ node1Config.getEngineConfig().getHttpConfig(),
+ node2Config.getEngineConfig().getHttpConfig());
+ Assertions.assertEquals(8080,
node2Config.getEngineConfig().getHttpConfig().getPort());
+ Assertions.assertNotEquals(
+ node1HttpPort,
+ node2HttpPort,
+ "node2 must expose the dynamically chosen REST port, not the
configured one");
+
+ // Embedded members share the Log4j context and job-log directory.
Both must have a
+ // log to list, otherwise counting distinct nodes could hide a fan-out
regression.
+ Path node1LogPath =
+ Paths.get(new
LogService(node1.node.getNodeEngine()).getLogPath()).toRealPath();
+ Path node2LogPath =
+ Paths.get(new
LogService(node2.node.getNodeEngine()).getLogPath()).toRealPath();
+ Assertions.assertEquals(node1LogPath, node2LogPath);
+ Assertions.assertTrue(
+ node1LogPath
+ .resolve("job-" + clientJobProxy.getJobId() + ".log")
+ .toFile()
+ .isFile());
+
+ List<Map<String, Object>> logEntries =
+ given().get(
+ buildHttpBaseUrl(node1HttpPort)
+ + RestConstant.REST_URL_LOGS
+ + "?format=JSON")
+ .then()
+ .statusCode(200)
+ .extract()
+ .as(new TypeRef<List<Map<String, Object>>>() {});
+
+ Set<String> reportedNodes =
+ logEntries.stream()
+ .map(entry -> String.valueOf(entry.get("node")))
+ .collect(Collectors.toSet());
+ Assertions.assertEquals(
+ new HashSet<>(expectedNodeIds()),
+ reportedNodes,
+ "GET /logs must enumerate both members, got " + reportedNodes);
+ Assertions.assertTrue(
+ reportedNodes.stream().anyMatch(node -> node.endsWith(":" +
node2HttpPort)),
+ "GET /logs must report node2 on its dynamic port, got " +
reportedNodes);
+ }
+
@Test
public void testLoggers() {
String loggersUrl =
HOST
- +
node1Config.getEngineConfig().getHttpConfig().getPort()
+ + httpPort(node1)
+
node1Config.getEngineConfig().getHttpConfig().getContextPath()
+ RestConstant.REST_URL_LOGGERS;
// a logger of this test only, so that changing its level cannot hide
job logs
@@ -328,7 +381,8 @@ public class RestApiIT {
.statusCode(200)
.body("scope", equalTo("cluster"))
.body("status", equalTo("SUCCESS"))
- .body("nodes", hasSize(ports.size()));
+ .body("nodes", hasSize(ports.size()))
+ .body("nodes.node",
containsInAnyOrder(expectedNodeIds().toArray(new String[0])));
given().post(loggersUrl + "/" + logger + "?scope=cluster&level=TRACE")
.then()
@@ -337,6 +391,7 @@ public class RestApiIT {
.body("status", equalTo("SUCCESS"))
.body("level", equalTo("TRACE"))
.body("nodes", hasSize(ports.size()))
+ .body("nodes.node",
containsInAnyOrder(expectedNodeIds().toArray(new String[0])))
.body("nodes[0].level", equalTo("TRACE"))
.body("nodes[0].origin", equalTo("runtime-override"));
@@ -345,9 +400,26 @@ public class RestApiIT {
.statusCode(200)
.body("scope", equalTo("cluster"))
.body("status", equalTo("SUCCESS"))
+ .body("nodes.node",
containsInAnyOrder(expectedNodeIds().toArray(new String[0])))
.body("nodes[0].origin", equalTo("file"));
}
+ private int httpPort(HazelcastInstanceImpl instance) {
+ SeaTunnelServer server =
+
instance.node.getNodeEngine().getService(SeaTunnelServer.SERVICE_NAME);
+ return server.getHttpPort();
+ }
+
+ private List<String> expectedNodeIds() {
+ return Arrays.asList(node1, node2).stream()
+ .map(
+ instance ->
+
instance.getCluster().getLocalMember().getAddress().getHost()
+ + ":"
+ + httpPort(instance))
+ .collect(Collectors.toList());
+ }
+
private CheckpointMonitorService resolveCheckpointMonitorService(
HazelcastInstanceImpl instance) {
try {
@@ -1350,8 +1422,7 @@ public class RestApiIT {
.until(
() -> {
Map<String, Object> overview =
- getCheckpointOverview(
- jobId,
buildHttpBaseUrl(httpPorts.get(0)));
+ getCheckpointOverview(jobId,
buildHttpBaseUrl(httpPort(node1)));
List<Map<String, Object>> pipelines =
castList(overview.get("pipelines"));
if (pipelines.isEmpty()) {
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/JettyService.java
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/JettyService.java
index a7866cad48..02d51290b5 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/JettyService.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/JettyService.java
@@ -103,34 +103,46 @@ public class JettyService {
private NodeEngineImpl nodeEngine;
private SeaTunnelConfig seaTunnelConfig;
+ private final ServerConnector httpConnector;
Server server;
public JettyService(NodeEngineImpl nodeEngine, SeaTunnelConfig
seaTunnelConfig) {
this.nodeEngine = nodeEngine;
this.seaTunnelConfig = seaTunnelConfig;
- int port = seaTunnelConfig.getEngineConfig().getHttpConfig().getPort();
- if
(seaTunnelConfig.getEngineConfig().getHttpConfig().isEnableDynamicPort()) {
- port =
- chooseAppropriatePort(
- port,
seaTunnelConfig.getEngineConfig().getHttpConfig().getPortRange());
- }
- log.info("SeaTunnel REST service will start on port {}", port);
+ HttpConfig httpConfig =
seaTunnelConfig.getEngineConfig().getHttpConfig();
this.server = new Server();
- if (seaTunnelConfig.getEngineConfig().getHttpConfig().isEnabled()) {
- // Enable http
- ServerConnector httpConnector = new ServerConnector(server);
+ if (httpConfig.isEnabled()) {
+ int port = httpConfig.getPort();
+ if (httpConfig.isEnableDynamicPort()) {
+ port = chooseAppropriatePort(port, httpConfig.getPortRange());
+ }
+ // LogService and LoggerLevelService must resolve each member's
own connector:
+ // multiple members can share the same HttpConfig.
+ httpConnector = new ServerConnector(server);
httpConnector.setPort(port);
server.addConnector(httpConnector);
+ } else {
+ httpConnector = null;
}
- if (seaTunnelConfig.getEngineConfig().getHttpConfig().isEnableHttps())
{
+ if (httpConfig.isEnableHttps()) {
// Enable https
- log.info("SeaTunnel REST service will start on https port {}",
port);
enableHttps(server, seaTunnelConfig);
}
}
+ /**
+ * Returns this node's bound HTTP port. Before binding, or when HTTP is
disabled, retains the
+ * configured-port fallback.
+ */
+ public int getHttpPort() {
+ int boundPort = httpConnector == null ? -1 :
httpConnector.getLocalPort();
+ return boundPort > 0
+ ? boundPort
+ : seaTunnelConfig.getEngineConfig().getHttpConfig().getPort();
+ }
+
public void enableHttps(Server server, SeaTunnelConfig seaTunnelConfig) {
HttpConfig httpConfig =
seaTunnelConfig.getEngineConfig().getHttpConfig();
@@ -278,6 +290,9 @@ public class JettyService {
try {
server.start();
+ if (httpConnector != null) {
+ log.info("SeaTunnel REST service started on http port {}",
getHttpPort());
+ }
} catch (Exception e) {
log.error("Jetty server start failed", e);
throw new RuntimeException(e);
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/SeaTunnelServer.java
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/SeaTunnelServer.java
index 37f39b03f8..de1837be1c 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/SeaTunnelServer.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/SeaTunnelServer.java
@@ -100,7 +100,7 @@ public class SeaTunnelServer
@Getter private CheckpointMonitorService checkpointMonitorService;
@Getter private ScheduledExecutorService monitorService;
private volatile RealtimeMetricsService realtimeMetricsService;
- private JettyService jettyService;
+ private volatile JettyService jettyService;
private TaskLogManagerService taskLogManagerService;
@Getter private SeaTunnelHealthMonitor seaTunnelHealthMonitor;
@@ -406,6 +406,14 @@ public class SeaTunnelServer
return seaTunnelConfig;
}
+ /** Returns this member's bound HTTP port, or the configured port before
Jetty is available. */
+ public int getHttpPort() {
+ JettyService service = jettyService;
+ return service == null
+ ? seaTunnelConfig.getEngineConfig().getHttpConfig().getPort()
+ : service.getHttpPort();
+ }
+
public NodeEngineImpl getNodeEngine() {
return nodeEngine;
}
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/operation/GetNodeHttpPortOperation.java
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/operation/GetNodeHttpPortOperation.java
index d5eed05f1a..b7a6f43ad9 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/operation/GetNodeHttpPortOperation.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/operation/GetNodeHttpPortOperation.java
@@ -24,7 +24,9 @@ import
com.hazelcast.nio.serialization.IdentifiedDataSerializable;
import com.hazelcast.spi.impl.AllowedDuringPassiveState;
import com.hazelcast.spi.impl.operationservice.Operation;
-/** Returns the REST HTTP port configured on the node that executes this
operation. */
+/**
+ * Returns the REST HTTP port bound by this member, with a configured-port
fallback before startup.
+ */
public class GetNodeHttpPortOperation extends Operation
implements IdentifiedDataSerializable, AllowedDuringPassiveState {
@@ -33,7 +35,7 @@ public class GetNodeHttpPortOperation extends Operation
@Override
public void run() {
SeaTunnelServer service = getService();
- response =
service.getSeaTunnelConfig().getEngineConfig().getHttpConfig().getPort();
+ response = service.getHttpPort();
}
@Override
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/LoggerLevelService.java
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/LoggerLevelService.java
index 628bc79336..e06dc641c2 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/LoggerLevelService.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/main/java/org/apache/seatunnel/engine/server/rest/service/LoggerLevelService.java
@@ -257,7 +257,11 @@ public class LoggerLevelService extends BaseService {
}
private String nodeId() {
- return nodeEngine.getThisAddress().getHost() + ":" +
httpConfig().getPort();
+ SeaTunnelServer seaTunnelServer = getSeaTunnelServer(false);
+ if (seaTunnelServer == null) {
+ throw new IllegalStateException("SeaTunnel server is not available
on this node.");
+ }
+ return nodeEngine.getThisAddress().getHost() + ":" +
seaTunnelServer.getHttpPort();
}
/**
diff --git
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/GetNodeHttpPortOperationTest.java
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/GetNodeHttpPortOperationTest.java
index 2db12c33db..9d7cdbd406 100644
---
a/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/GetNodeHttpPortOperationTest.java
+++
b/seatunnel-engine/seatunnel-engine-server/src/test/java/org/apache/seatunnel/engine/server/operation/GetNodeHttpPortOperationTest.java
@@ -18,28 +18,52 @@
package org.apache.seatunnel.engine.server.operation;
import org.apache.seatunnel.engine.common.config.SeaTunnelConfig;
+import org.apache.seatunnel.engine.common.config.server.HttpConfig;
import org.apache.seatunnel.engine.server.AbstractSeaTunnelServerTest;
+import org.apache.seatunnel.engine.server.JettyService;
+import org.apache.seatunnel.engine.server.SeaTunnelServer;
+import org.apache.seatunnel.engine.server.TestUtils;
import org.apache.seatunnel.engine.server.utils.NodeEngineUtil;
import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.Test;
import com.hazelcast.cluster.Address;
+import java.io.IOException;
+import java.io.UncheckedIOException;
+import java.net.ServerSocket;
+import java.net.URL;
+
public class GetNodeHttpPortOperationTest
extends AbstractSeaTunnelServerTest<GetNodeHttpPortOperationTest> {
- private static final int HTTP_PORT = 18085;
+ private static final int HTTP_PORT = TestUtils.getAvailablePort(100);
+
+ @BeforeAll
+ @Override
+ public void before() {
+ // Keep the configured port occupied until Jetty has selected and
bound another port.
+ try (ServerSocket occupied = new ServerSocket(HTTP_PORT)) {
+ super.before();
+ } catch (IOException e) {
+ throw new UncheckedIOException(e);
+ }
+ }
@Override
public SeaTunnelConfig loadSeaTunnelConfig() {
SeaTunnelConfig config = super.loadSeaTunnelConfig();
+ config.getEngineConfig().setHttpConfig(new HttpConfig());
config.getEngineConfig().getHttpConfig().setPort(HTTP_PORT);
+ config.getEngineConfig().getHttpConfig().setEnabled(true);
+ config.getEngineConfig().getHttpConfig().setEnableDynamicPort(true);
return config;
}
@Test
- public void testReturnsConfiguredHttpPort() throws Exception {
+ public void testReturnsBoundHttpPortWithoutChangingConfiguration() throws
Exception {
Address localAddress =
instance.getCluster().getLocalMember().getAddress();
int result =
@@ -48,6 +72,44 @@ public class GetNodeHttpPortOperationTest
nodeEngine, new
GetNodeHttpPortOperation(), localAddress)
.get();
- Assertions.assertEquals(HTTP_PORT, result);
+ Assertions.assertNotEquals(
+ HTTP_PORT, result, "The dynamically bound port must be
published to peers");
+ Assertions.assertEquals(server.getHttpPort(), result);
+ Assertions.assertEquals(
+ HTTP_PORT,
server.getSeaTunnelConfig().getEngineConfig().getHttpConfig().getPort());
+ }
+
+ @Test
+ public void testConfiguredPortFallbackBeforeStartup() {
+ SeaTunnelConfig config = new SeaTunnelConfig();
+ config.getEngineConfig().setHttpConfig(new HttpConfig());
+ config.getEngineConfig().getHttpConfig().setPort(18085);
+ Assertions.assertEquals(18085, new
SeaTunnelServer(config).getHttpPort());
+ }
+
+ @Test
+ public void testHttpsOnlyDoesNotProbeTheDisabledHttpPort() throws
Exception {
+ try (ServerSocket occupied = new ServerSocket(0)) {
+ SeaTunnelConfig config = new SeaTunnelConfig();
+ config.getEngineConfig().setHttpConfig(new HttpConfig());
+
config.getEngineConfig().getHttpConfig().setPort(occupied.getLocalPort());
+ config.getEngineConfig().getHttpConfig().setPortRange(0);
+
config.getEngineConfig().getHttpConfig().setEnableDynamicPort(true);
+ config.getEngineConfig().getHttpConfig().setEnabled(false);
+ config.getEngineConfig().getHttpConfig().setEnableHttps(true);
+ URL keyStore =
getClass().getClassLoader().getResource("https/server_keystore.jks");
+ Assertions.assertNotNull(keyStore);
+
config.getEngineConfig().getHttpConfig().setKeyStorePath(keyStore.toExternalForm());
+ // An omitted password makes Jetty read from Surefire's command
stream on stdin.
+ config.getEngineConfig()
+ .getHttpConfig()
+ .setKeyStorePassword("server_keystore_password");
+ config.getEngineConfig()
+ .getHttpConfig()
+ .setKeyManagerPassword("server_keystore_password");
+
+ JettyService service = new
JettyService(instance.node.getNodeEngine(), config);
+ Assertions.assertEquals(occupied.getLocalPort(),
service.getHttpPort());
+ }
}
}