This is an automated email from the ASF dual-hosted git repository.
dockerzhang 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 0709547f2e [INLONG-9606][Manager] Fix the problem of incorrect flow
status when cls sink configuration fails (#9607)
0709547f2e is described below
commit 0709547f2e630bc0990a98bf061a045ddf2fac05
Author: fuweng11 <[email protected]>
AuthorDate: Tue Jan 23 12:36:42 2024 +0800
[INLONG-9606][Manager] Fix the problem of incorrect flow status when cls
sink configuration fails (#9607)
---
.../manager/service/node/cls/ClsDataNodeOperator.java | 3 +--
.../manager/service/resource/sink/cls/ClsOperator.java | 13 ++++++-------
.../service/resource/sink/cls/ClsResourceOperator.java | 5 ++---
.../manager/web/controller/InlongClusterController.java | 3 ++-
4 files changed, 11 insertions(+), 13 deletions(-)
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 40ec26f210..f290b6c81a 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
@@ -33,7 +33,6 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import com.tencentcloudapi.cls.v20201016.ClsClient;
import com.tencentcloudapi.cls.v20201016.models.DescribeTopicsRequest;
import com.tencentcloudapi.common.Credential;
-import com.tencentcloudapi.common.exception.TencentCloudSDKException;
import com.tencentcloudapi.common.profile.ClientProfile;
import com.tencentcloudapi.common.profile.HttpProfile;
import org.apache.commons.lang3.StringUtils;
@@ -99,7 +98,7 @@ public class ClsDataNodeOperator extends
AbstractDataNodeOperator {
DescribeTopicsRequest req = new DescribeTopicsRequest();
try {
client.DescribeTopics(req);
- } catch (TencentCloudSDKException e) {
+ } catch (Exception e) {
String errMsg = String.format("connect tencent cloud error
endPoint = %s secretId = %s secretKey = %s",
dataNodeRequest.getEndpoint(),
dataNodeRequest.getManageSecretId(),
dataNodeRequest.getManageSecretKey());
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsOperator.java
index d15223b861..f5162d1cc0 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsOperator.java
@@ -39,7 +39,6 @@ import com.tencentcloudapi.cls.v20201016.models.RuleInfo;
import com.tencentcloudapi.cls.v20201016.models.Tag;
import com.tencentcloudapi.cls.v20201016.models.TopicInfo;
import com.tencentcloudapi.common.Credential;
-import com.tencentcloudapi.common.exception.TencentCloudSDKException;
import com.tencentcloudapi.common.profile.ClientProfile;
import com.tencentcloudapi.common.profile.HttpProfile;
import org.apache.commons.lang3.ArrayUtils;
@@ -65,7 +64,7 @@ public class ClsOperator {
public String createTopicReturnTopicId(String topicName, String logSetId,
String tag, Integer storageDuration,
String secretId, String secretKey, String region)
- throws TencentCloudSDKException {
+ throws Exception {
ClsClient client = getClsClient(secretId, secretKey, region);
CreateTopicRequest req = getCreateTopicRequest(tag, logSetId,
topicName, storageDuration);
CreateTopicResponse resp = client.CreateTopic(req);
@@ -76,7 +75,7 @@ public class ClsOperator {
}
public void updateTopicTag(String topicId, String tag, String secretId,
String secretKey, String region)
- throws TencentCloudSDKException {
+ throws Exception {
ClsClient client = getClsClient(secretId, secretKey, region);
ModifyTopicRequest modifyTopicRequest = new ModifyTopicRequest();
modifyTopicRequest.setTags(convertTags(tag.split(InlongConstants.CENTER_LINE)));
@@ -109,7 +108,7 @@ public class ClsOperator {
CreateIndexResponse createIndexResponse =
clsClient.CreateIndex(req);
LOG.debug("create index success for topic = {}, tokenizer = {},
requestId = {}", topicId,
tokenizer, createIndexResponse.getRequestId());
- } catch (TencentCloudSDKException e) {
+ } catch (Exception e) {
String errMsg = "Create cls topic index failed: " + e.getMessage();
LOG.error(errMsg, e);
throw new BusinessException(errMsg);
@@ -134,7 +133,7 @@ public class ClsOperator {
return topics[0].getTopicId();
}
return null;
- } catch (TencentCloudSDKException e) {
+ } catch (Exception e) {
String errMsg = "describe cls topic failed: " + e.getMessage();
LOG.error(errMsg, e);
throw new BusinessException(errMsg);
@@ -165,7 +164,7 @@ public class ClsOperator {
try {
DescribeIndexResponse resp = clsClient.DescribeIndex(req);
return resp.getRule() == null ? null :
resp.getRule().getFullText();
- } catch (TencentCloudSDKException e) {
+ } catch (Exception e) {
String errMsg = "describe cls topic index failed: " +
e.getMessage();
LOG.error(errMsg, e);
throw new BusinessException(errMsg);
@@ -186,7 +185,7 @@ public class ClsOperator {
ModifyIndexResponse modifyIndexResponse =
clsClient.ModifyIndex(req);
LOG.debug("update index success for topicId = {}, tokenizer = {},
requestId = {}", topicId, tokenizer,
modifyIndexResponse.getRequestId());
- } catch (TencentCloudSDKException e) {
+ } catch (Exception e) {
String errMsg = "update cls topic index failed: " + e.getMessage();
LOG.error(errMsg, e);
throw new BusinessException(errMsg);
diff --git
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsResourceOperator.java
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsResourceOperator.java
index 173a139758..7e2d297466 100644
---
a/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsResourceOperator.java
+++
b/inlong-manager/manager-service/src/main/java/org/apache/inlong/manager/service/resource/sink/cls/ClsResourceOperator.java
@@ -33,7 +33,6 @@ import org.apache.inlong.manager.pojo.sink.cls.ClsSinkDTO;
import
org.apache.inlong.manager.service.resource.sink.AbstractStandaloneSinkResourceOperator;
import org.apache.inlong.manager.service.sink.StreamSinkService;
-import com.tencentcloudapi.common.exception.TencentCloudSDKException;
import org.apache.commons.lang3.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -93,7 +92,7 @@ public class ClsResourceOperator extends
AbstractStandaloneSinkResourceOperator
sinkService.updateStatus(sinkInfo.getId(),
SinkStatus.CONFIG_SUCCESSFUL.getCode(), info);
LOG.info("update cls info status success for sinkId= {}, topicName
= {}", sinkInfo.getSinkName(),
clsSinkDTO.getTopicName());
- } catch (TencentCloudSDKException e) {
+ } catch (Exception e) {
String errMsg = "Create cls topic failed: " + e.getMessage();
LOG.error(errMsg, e);
sinkService.updateStatus(sinkInfo.getId(),
SinkStatus.CONFIG_FAILED.getCode(), errMsg);
@@ -102,7 +101,7 @@ public class ClsResourceOperator extends
AbstractStandaloneSinkResourceOperator
}
private String getTopicID(ClsDataNodeDTO clsDataNode, ClsSinkDTO
clsSinkDTO)
- throws TencentCloudSDKException {
+ throws Exception {
String topicId =
clsOperator.describeTopicIDByTopicName(clsSinkDTO.getTopicName(),
clsDataNode.getLogSetId(),
clsDataNode.getManageSecretId(),
clsDataNode.getManageSecretKey(),
clsDataNode.getRegion());
diff --git
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongClusterController.java
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongClusterController.java
index bba545a3df..f0483bf78c 100644
---
a/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongClusterController.java
+++
b/inlong-manager/manager-web/src/main/java/org/apache/inlong/manager/web/controller/InlongClusterController.java
@@ -47,6 +47,7 @@ import io.swagger.annotations.Api;
import io.swagger.annotations.ApiImplicitParam;
import io.swagger.annotations.ApiImplicitParams;
import io.swagger.annotations.ApiOperation;
+import org.apache.shiro.authz.annotation.Logical;
import org.apache.shiro.authz.annotation.RequiresRoles;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.validation.annotation.Validated;
@@ -158,7 +159,7 @@ public class InlongClusterController {
@PostMapping(value = "/cluster/save")
@ApiOperation(value = "Save cluster")
@OperationLog(operation = OperationType.CREATE, operationTarget =
OperationTarget.CLUSTER)
- @RequiresRoles(value = UserRoleCode.TENANT_ADMIN)
+ @RequiresRoles(logical = Logical.OR, value = {UserRoleCode.INLONG_ADMIN,
UserRoleCode.TENANT_ADMIN})
public Response<Integer> save(@Validated(SaveValidation.class)
@RequestBody ClusterRequest request) {
String currentUser = LoginUserUtils.getLoginUser().getName();
return Response.success(clusterService.save(request, currentUser));