This is an automated email from the ASF dual-hosted git repository.

fuweng11 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new 78636fd312 [INLONG-12193][Manager] Fix SSRF in DataNode testConnection 
(#12194)
78636fd312 is described below

commit 78636fd31210cebeb6d7d47c14b68156ed0e62ca
Author: fuweng11 <[email protected]>
AuthorDate: Thu Aug 27 10:55:29 2026 +0800

    [INLONG-12193][Manager] Fix SSRF in DataNode testConnection (#12194)
    
    Co-authored-by: wakefu <[email protected]>
---
 .../manager/common/util/UrlVerificationUtils.java  | 19 ++++++++++++++
 .../service/node/cls/ClsDataNodeOperator.java      | 13 +++++++++-
 .../service/node/hudi/HudiDataNodeOperator.java    |  8 ++++++
 .../node/iceberg/IcebergDataNodeOperator.java      |  9 +++++++
 .../service/node/kafka/KafkaDataNodeOperator.java  |  8 ++++++
 .../service/node/kudu/KuduDataNodeOperator.java    |  8 ++++++
 .../node/pulsar/PulsarDataNodeOperator.java        |  8 ++++++
 .../service/node/redis/RedisDataNodeOperator.java  | 30 ++++++++++++++++++++++
 8 files changed, 102 insertions(+), 1 deletion(-)

diff --git 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/util/UrlVerificationUtils.java
 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/util/UrlVerificationUtils.java
index e7f3814c02..2af04e38e8 100644
--- 
a/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/util/UrlVerificationUtils.java
+++ 
b/inlong-manager/manager-common/src/main/java/org/apache/inlong/manager/common/util/UrlVerificationUtils.java
@@ -102,6 +102,25 @@ public class UrlVerificationUtils {
         validateHostNotInternal(host);
     }
 
+    public static void validateEndpointListNotInternal(String endpoints) 
throws Exception {
+        if (endpoints == null || endpoints.trim().isEmpty()) {
+            throw new Exception("endpoint list cannot be null or empty");
+        }
+        String[] parts = endpoints.split("[,;]");
+        boolean hasEntry = false;
+        for (String part : parts) {
+            String entry = part == null ? null : part.trim();
+            if (entry == null || entry.isEmpty()) {
+                continue;
+            }
+            hasEntry = true;
+            validateUrlNotInternal(entry);
+        }
+        if (!hasEntry) {
+            throw new Exception("endpoint list cannot be blank");
+        }
+    }
+
     /**
      * Validates that a hostname does not resolve to an internal/private IP 
address.
      *
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/cls/ClsDataNodeOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/cls/ClsDataNodeOperator.java
index 65b7a17b25..730617ac02 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/cls/ClsDataNodeOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/cls/ClsDataNodeOperator.java
@@ -23,6 +23,7 @@ import org.apache.inlong.manager.common.consts.DataNodeType;
 import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
 import org.apache.inlong.manager.common.exceptions.BusinessException;
 import org.apache.inlong.manager.common.util.CommonBeanUtils;
+import org.apache.inlong.manager.common.util.UrlVerificationUtils;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.pojo.node.DataNodeInfo;
 import org.apache.inlong.manager.pojo.node.DataNodeRequest;
@@ -91,9 +92,19 @@ public class ClsDataNodeOperator extends 
AbstractDataNodeOperator {
     @Override
     public Boolean testConnection(DataNodeRequest request) {
         ClsDataNodeRequest dataNodeRequest = (ClsDataNodeRequest) request;
+        // SSRF protection: the endpoint used to contact Tencent Cloud must be 
an external address
+        String endpoint = dataNodeRequest.getEndpoint();
+        if (StringUtils.isNotBlank(endpoint)) {
+            try {
+                UrlVerificationUtils.validateUrlNotInternal(endpoint);
+            } catch (Exception e) {
+                throw new BusinessException(ErrorCodeEnum.INVALID_PARAMETER,
+                        "SSRF protection: " + e.getMessage());
+            }
+        }
         Credential cred = new Credential(dataNodeRequest.getManageSecretId(), 
dataNodeRequest.getManageSecretId());
         HttpProfile httpProfile = new HttpProfile();
-        httpProfile.setEndpoint(dataNodeRequest.getEndpoint());
+        httpProfile.setEndpoint(endpoint);
         ClientProfile clientProfile = new ClientProfile();
         clientProfile.setHttpProfile(httpProfile);
         ClsClient client = new ClsClient(cred, dataNodeRequest.getRegion(), 
clientProfile);
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/hudi/HudiDataNodeOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/hudi/HudiDataNodeOperator.java
index 7f9928eb66..e320ce3861 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/hudi/HudiDataNodeOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/hudi/HudiDataNodeOperator.java
@@ -22,6 +22,7 @@ import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
 import org.apache.inlong.manager.common.exceptions.BusinessException;
 import org.apache.inlong.manager.common.util.CommonBeanUtils;
 import org.apache.inlong.manager.common.util.Preconditions;
+import org.apache.inlong.manager.common.util.UrlVerificationUtils;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.pojo.node.DataNodeInfo;
 import org.apache.inlong.manager.pojo.node.DataNodeRequest;
@@ -90,6 +91,13 @@ public class HudiDataNodeOperator extends 
AbstractDataNodeOperator {
         String metastoreUri = hudiRequest.getUrl();
         String warehouse = hudiRequest.getWarehouse();
         Preconditions.expectNotBlank(metastoreUri, 
ErrorCodeEnum.INVALID_PARAMETER, "connection url cannot be empty");
+        // SSRF protection: block requests to internal/loopback/cloud-metadata 
addresses.
+        try {
+            UrlVerificationUtils.validateUrlNotInternal(metastoreUri);
+        } catch (Exception e) {
+            throw new BusinessException(ErrorCodeEnum.INVALID_PARAMETER,
+                    "SSRF protection: " + e.getMessage());
+        }
         try (HudiCatalogClient client = new HudiCatalogClient(metastoreUri, 
warehouse)) {
             client.open();
             client.listAllDatabases();
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/iceberg/IcebergDataNodeOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/iceberg/IcebergDataNodeOperator.java
index 809445c80f..46c653e6a0 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/iceberg/IcebergDataNodeOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/iceberg/IcebergDataNodeOperator.java
@@ -22,6 +22,7 @@ import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
 import org.apache.inlong.manager.common.exceptions.BusinessException;
 import org.apache.inlong.manager.common.util.CommonBeanUtils;
 import org.apache.inlong.manager.common.util.Preconditions;
+import org.apache.inlong.manager.common.util.UrlVerificationUtils;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.pojo.node.DataNodeInfo;
 import org.apache.inlong.manager.pojo.node.DataNodeRequest;
@@ -90,6 +91,14 @@ public class IcebergDataNodeOperator extends 
AbstractDataNodeOperator {
         String metastoreUri = icebergDataNodeRequest.getUrl();
         String warehouse = icebergDataNodeRequest.getWarehouse();
         Preconditions.expectNotBlank(metastoreUri, 
ErrorCodeEnum.INVALID_PARAMETER, "connection url cannot be empty");
+        // SSRF protection: even though the service layer validates 
request.getUrl(), keep an
+        // operator-side check so the operator remains safe if invoked from 
other call sites.
+        try {
+            UrlVerificationUtils.validateUrlNotInternal(metastoreUri);
+        } catch (Exception e) {
+            throw new BusinessException(ErrorCodeEnum.INVALID_PARAMETER,
+                    "SSRF protection: " + e.getMessage());
+        }
         try {
             HiveCatalog catalog = IcebergCatalogUtils.getCatalog(metastoreUri, 
warehouse);
             catalog.listNamespaces();
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/kafka/KafkaDataNodeOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/kafka/KafkaDataNodeOperator.java
index bf57e85fe5..9796d00ca8 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/kafka/KafkaDataNodeOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/kafka/KafkaDataNodeOperator.java
@@ -24,6 +24,7 @@ import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
 import org.apache.inlong.manager.common.exceptions.BusinessException;
 import org.apache.inlong.manager.common.util.CommonBeanUtils;
 import org.apache.inlong.manager.common.util.Preconditions;
+import org.apache.inlong.manager.common.util.UrlVerificationUtils;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.pojo.cluster.kafka.KafkaClusterInfo;
 import org.apache.inlong.manager.pojo.node.DataNodeInfo;
@@ -114,6 +115,13 @@ public class KafkaDataNodeOperator extends 
AbstractDataNodeOperator {
         String bootstrapServers = kafkaDataNodeRequest.getBootstrapServers();
         Preconditions.expectNotBlank(bootstrapServers, 
ErrorCodeEnum.INVALID_PARAMETER,
                 "connection bootstrapServers  cannot be empty");
+        // SSRF protection: every bootstrap server must resolve to a 
public/external address
+        try {
+            
UrlVerificationUtils.validateEndpointListNotInternal(bootstrapServers);
+        } catch (Exception e) {
+            throw new BusinessException(ErrorCodeEnum.INVALID_PARAMETER,
+                    "SSRF protection: " + e.getMessage());
+        }
         if (getKafkaConnection(bootstrapServers)) {
             LOGGER.info("kafka connection success for bootstrapServers={}",
                     bootstrapServers);
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/kudu/KuduDataNodeOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/kudu/KuduDataNodeOperator.java
index 6d1b11c4aa..dd8a9da2df 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/kudu/KuduDataNodeOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/kudu/KuduDataNodeOperator.java
@@ -23,6 +23,7 @@ import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
 import org.apache.inlong.manager.common.exceptions.BusinessException;
 import org.apache.inlong.manager.common.util.CommonBeanUtils;
 import org.apache.inlong.manager.common.util.Preconditions;
+import org.apache.inlong.manager.common.util.UrlVerificationUtils;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.pojo.node.DataNodeInfo;
 import org.apache.inlong.manager.pojo.node.DataNodeRequest;
@@ -96,6 +97,13 @@ public class KuduDataNodeOperator extends 
AbstractDataNodeOperator {
         KuduDataNodeRequest kuduRequest = (KuduDataNodeRequest) request;
         String masters = kuduRequest.getMasters();
         Preconditions.expectNotBlank(masters, ErrorCodeEnum.INVALID_PARAMETER, 
"masters cannot be empty");
+        // SSRF protection: every master endpoint must resolve to a 
public/external address
+        try {
+            UrlVerificationUtils.validateEndpointListNotInternal(masters);
+        } catch (Exception e) {
+            throw new BusinessException(ErrorCodeEnum.INVALID_PARAMETER,
+                    "SSRF protection: " + e.getMessage());
+        }
 
         try (KuduResourceClient kuduClient = new KuduResourceClient(masters)) {
             kuduClient.getTablesList();
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/pulsar/PulsarDataNodeOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/pulsar/PulsarDataNodeOperator.java
index 26a7fa0a55..5cae9d654d 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/pulsar/PulsarDataNodeOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/pulsar/PulsarDataNodeOperator.java
@@ -24,6 +24,7 @@ import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
 import org.apache.inlong.manager.common.exceptions.BusinessException;
 import org.apache.inlong.manager.common.util.CommonBeanUtils;
 import org.apache.inlong.manager.common.util.Preconditions;
+import org.apache.inlong.manager.common.util.UrlVerificationUtils;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.pojo.cluster.pulsar.PulsarClusterInfo;
 import org.apache.inlong.manager.pojo.node.DataNodeInfo;
@@ -100,6 +101,13 @@ public class PulsarDataNodeOperator extends 
AbstractDataNodeOperator {
         String adminUrl = pulsarDataNodeRequest.getAdminUrl();
         String token = pulsarDataNodeRequest.getToken();
         Preconditions.expectNotBlank(adminUrl, 
ErrorCodeEnum.INVALID_PARAMETER, "connection admin urlcannot be empty");
+        // SSRF protection: block requests to 
internal/loopback/link-local/cloud-metadata addresses
+        try {
+            UrlVerificationUtils.validateUrlNotInternal(adminUrl);
+        } catch (Exception e) {
+            throw new BusinessException(ErrorCodeEnum.INVALID_PARAMETER,
+                    "SSRF protection: " + e.getMessage());
+        }
         if (getPulsarConnection(adminUrl, token)) {
             LOGGER.info("pulsar  connection success for adminUrl={}, token={}",
                     adminUrl, token);
diff --git 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/redis/RedisDataNodeOperator.java
 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/redis/RedisDataNodeOperator.java
index 46d60e91f1..7bc743b73a 100644
--- 
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/redis/RedisDataNodeOperator.java
+++ 
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/node/redis/RedisDataNodeOperator.java
@@ -22,6 +22,7 @@ import org.apache.inlong.manager.common.enums.ErrorCodeEnum;
 import org.apache.inlong.manager.common.exceptions.BusinessException;
 import org.apache.inlong.manager.common.util.CommonBeanUtils;
 import org.apache.inlong.manager.common.util.Preconditions;
+import org.apache.inlong.manager.common.util.UrlVerificationUtils;
 import org.apache.inlong.manager.dao.entity.DataNodeEntity;
 import org.apache.inlong.manager.pojo.node.DataNodeInfo;
 import org.apache.inlong.manager.pojo.node.DataNodeRequest;
@@ -117,6 +118,35 @@ public class RedisDataNodeOperator extends 
AbstractDataNodeOperator {
     @Override
     public Boolean testConnection(DataNodeRequest request) {
         RedisDataNodeRequest redisDataNodeRequest = (RedisDataNodeRequest) 
request;
+        // SSRF protection: reject requests targeting 
internal/loopback/cloud-metadata addresses.
+        // Redis operator uses different fields per cluster mode; validate all 
that apply.
+        try {
+            RedisClusterMode clusterMode = 
RedisClusterMode.of(redisDataNodeRequest.getClusterMode());
+            switch (clusterMode) {
+                case STANDALONE:
+                    if 
(StringUtils.isNotBlank(redisDataNodeRequest.getHost())) {
+                        
UrlVerificationUtils.validateHostNotInternal(redisDataNodeRequest.getHost());
+                    }
+                    break;
+                case CLUSTER:
+                    if 
(StringUtils.isNotBlank(redisDataNodeRequest.getClusterNodes())) {
+                        UrlVerificationUtils
+                                
.validateEndpointListNotInternal(redisDataNodeRequest.getClusterNodes());
+                    }
+                    break;
+                case SENTINEL:
+                    if 
(StringUtils.isNotBlank(redisDataNodeRequest.getSentinelsInfo())) {
+                        UrlVerificationUtils
+                                
.validateEndpointListNotInternal(redisDataNodeRequest.getSentinelsInfo());
+                    }
+                    break;
+                default:
+                    break;
+            }
+        } catch (Exception e) {
+            throw new BusinessException(ErrorCodeEnum.INVALID_PARAMETER,
+                    "SSRF protection: " + e.getMessage());
+        }
         try {
             return RedisResourceClient.testConnection(redisDataNodeRequest);
         } catch (Exception e) {

Reply via email to