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());
+        }
     }
 }

Reply via email to