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/incubator-inlong.git


The following commit(s) were added to refs/heads/master by this push:
     new d9162c2d0 [INLONG-3621][TubeMQ] Added cluster switching method and 
delete cluster support to delete master (#3630)
d9162c2d0 is described below

commit d9162c2d0f748128fe4fe91b4d221cd6e81aa255
Author: bluewang <[email protected]>
AuthorDate: Tue Apr 12 17:41:28 2022 +0800

    [INLONG-3621][TubeMQ] Added cluster switching method and delete cluster 
support to delete master (#3630)
---
 .../controller/cluster/ClusterController.java      |  3 +++
 .../cluster/request/SwitchClusterReq.java          | 30 ++++++++++++++++++++++
 .../manager/repository/MasterRepository.java       |  7 +++++
 .../tubemq/manager/service/ClusterServiceImpl.java |  5 ++++
 .../tubemq/manager/service/MasterServiceImpl.java  |  9 +++++++
 .../inlong/tubemq/manager/service/TubeConst.java   |  1 +
 .../manager/service/interfaces/MasterService.java  |  8 ++++++
 7 files changed, 63 insertions(+)

diff --git 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/ClusterController.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/ClusterController.java
index 5530e4691..99c1903e0 100644
--- 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/ClusterController.java
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/ClusterController.java
@@ -29,6 +29,7 @@ import 
org.apache.inlong.tubemq.manager.controller.TubeMQResult;
 import org.apache.inlong.tubemq.manager.controller.cluster.dto.ClusterDto;
 import 
org.apache.inlong.tubemq.manager.controller.cluster.request.AddClusterReq;
 import 
org.apache.inlong.tubemq.manager.controller.cluster.request.DeleteClusterReq;
+import 
org.apache.inlong.tubemq.manager.controller.cluster.request.SwitchClusterReq;
 import org.apache.inlong.tubemq.manager.controller.cluster.vo.ClusterVo;
 import 
org.apache.inlong.tubemq.manager.controller.group.result.ConsumerGroupInfoRes;
 import 
org.apache.inlong.tubemq.manager.controller.group.result.ConsumerInfoRes;
@@ -81,6 +82,8 @@ public class ClusterController {
                 return deleteCluster(gson.fromJson(req, 
DeleteClusterReq.class));
             case TubeConst.MODIFY:
                 return changeCluster(gson.fromJson(req, ClusterDto.class));
+            case TubeConst.SWITCH:
+                return masterService.baseRequestMaster(gson.fromJson(req, 
SwitchClusterReq.class));
             default:
                 return 
TubeMQResult.errorResult(TubeMQErrorConst.NO_SUCH_METHOD);
         }
diff --git 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/request/SwitchClusterReq.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/request/SwitchClusterReq.java
new file mode 100644
index 000000000..d03cc5ef6
--- /dev/null
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/controller/cluster/request/SwitchClusterReq.java
@@ -0,0 +1,30 @@
+/*
+ * 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.inlong.tubemq.manager.controller.cluster.request;
+
+import lombok.Data;
+import lombok.EqualsAndHashCode;
+import lombok.ToString;
+import org.apache.inlong.tubemq.manager.controller.node.request.BaseReq;
+
+@Data
+@EqualsAndHashCode(callSuper = true)
+@ToString(callSuper = true)
+public class SwitchClusterReq extends BaseReq {
+    private String confModAuthToken;
+}
diff --git 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/repository/MasterRepository.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/repository/MasterRepository.java
index e0b653664..54d538115 100644
--- 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/repository/MasterRepository.java
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/repository/MasterRepository.java
@@ -56,4 +56,11 @@ public interface MasterRepository extends 
JpaRepository<MasterEntry, Long> {
      * @return
      */
     List<MasterEntry> findMasterEntryByIpEquals(String masterIp);
+
+    /**
+     * delete master by cluster id
+     *
+     * @return
+     */
+    Integer deleteByClusterId(Long clusterId);
 }
diff --git 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/ClusterServiceImpl.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/ClusterServiceImpl.java
index 4fcd2963b..151367c8e 100644
--- 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/ClusterServiceImpl.java
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/ClusterServiceImpl.java
@@ -31,6 +31,7 @@ import org.apache.inlong.tubemq.manager.entry.ClusterEntry;
 import org.apache.inlong.tubemq.manager.entry.MasterEntry;
 import org.apache.inlong.tubemq.manager.repository.ClusterRepository;
 import org.apache.inlong.tubemq.manager.service.interfaces.ClusterService;
+import org.apache.inlong.tubemq.manager.service.interfaces.MasterService;
 import org.apache.inlong.tubemq.manager.service.interfaces.NodeService;
 import org.springframework.beans.factory.annotation.Autowired;
 import org.springframework.stereotype.Component;
@@ -45,6 +46,9 @@ public class ClusterServiceImpl implements ClusterService {
     @Autowired
     NodeService nodeService;
 
+    @Autowired
+    MasterService masterService;
+
     @Override
     @Transactional(rollbackOn = Exception.class)
     public void addClusterAndMasterNode(AddClusterReq req) {
@@ -64,6 +68,7 @@ public class ClusterServiceImpl implements ClusterService {
     @Override
     @Transactional(rollbackOn = Exception.class)
     public void deleteCluster(Long clusterId) {
+        masterService.deleteMaster(clusterId);
         Integer successCode = clusterRepository.deleteByClusterId(clusterId);
         if (successCode.equals(DELETE_FAIL)) {
             throw new RuntimeException("no such cluster with clusterId = " + 
clusterId);
diff --git 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/MasterServiceImpl.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/MasterServiceImpl.java
index 3d6379544..4d24133b7 100644
--- 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/MasterServiceImpl.java
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/MasterServiceImpl.java
@@ -37,6 +37,7 @@ import 
org.apache.inlong.tubemq.manager.controller.TubeMQResult;
 import org.apache.inlong.tubemq.manager.controller.node.request.BaseReq;
 import org.apache.inlong.tubemq.manager.entry.MasterEntry;
 import org.apache.inlong.tubemq.manager.repository.MasterRepository;
+import static org.apache.inlong.tubemq.manager.service.TubeConst.DELETE_FAIL;
 import org.apache.inlong.tubemq.manager.service.interfaces.MasterService;
 import org.apache.inlong.tubemq.manager.service.tube.TubeHttpResponse;
 import org.apache.inlong.tubemq.manager.utils.ConvertUtils;
@@ -189,4 +190,12 @@ public class MasterServiceImpl implements MasterService {
                 + method + "&" + "clusterId=" + clusterId;
     }
 
+    @Override
+    public void deleteMaster(Long clusterId) {
+        Integer successCode = masterRepository.deleteByClusterId(clusterId);
+        if (successCode.equals(DELETE_FAIL)) {
+            throw new RuntimeException("no such master with clusterId = " + 
clusterId);
+        }
+    }
+
 }
diff --git 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/TubeConst.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/TubeConst.java
index ef851f9b8..87a1205ac 100644
--- 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/TubeConst.java
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/TubeConst.java
@@ -51,6 +51,7 @@ public class TubeConst {
     public static final String CLONE = "clone";
     public static final String ADD = "add";
     public static final String QUERY = "query";
+    public static final String SWITCH = "switch";
     public static final String REBALANCE_CONSUMER_GROUP = "rebalanceGroup";
     public static final String REBALANCE_CONSUMER = "rebalanceConsumer";
     public static final String SET_READ_OR_WRITE = "setReadOrWrite";
diff --git 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/interfaces/MasterService.java
 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/interfaces/MasterService.java
index 0ac583b24..3e062fc95 100644
--- 
a/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/interfaces/MasterService.java
+++ 
b/inlong-tubemq/tubemq-manager/src/main/java/org/apache/inlong/tubemq/manager/service/interfaces/MasterService.java
@@ -101,4 +101,12 @@ public interface MasterService {
     TubeMQResult checkMasterNodeStatus(String masterIp, Integer masterPort);
 
     String getQueryCountUrl(Integer clusterId, String method);
+
+    /**
+     * delete master by cluster id
+     *
+     * @param clusterId
+     */
+    void deleteMaster(Long clusterId);
+
 }

Reply via email to