mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-01 05:46:12 +08:00
invert dependency between remoting and common module (#5503)
* invert dependency between remoting and common module * fix bazel * Fix Bazel warning of indirect dependency * remove duplicate class ServiceThread * fix conflict * revert delete class ServiceThread because of introduced by dledger Co-authored-by: Li Zhanhui <lizhanhui@gmail.com>
This commit is contained in:
@@ -19,14 +19,14 @@ package org.apache.rocketmq.controller;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import org.apache.rocketmq.common.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.CleanControllerBrokerDataRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.remoting.RemotingServer;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.remoting.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.CleanControllerBrokerDataRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
|
||||
/**
|
||||
* The api for controller
|
||||
|
||||
@@ -24,31 +24,30 @@ import java.util.concurrent.LinkedBlockingQueue;
|
||||
import java.util.concurrent.RunnableFuture;
|
||||
import java.util.concurrent.ThreadPoolExecutor;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.Configuration;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.ThreadFactoryImpl;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.future.FutureTaskExt;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.protocol.RequestCode;
|
||||
import org.apache.rocketmq.common.protocol.body.BrokerMemberGroup;
|
||||
import org.apache.rocketmq.common.protocol.header.NotifyBrokerRoleChangedRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterResponseHeader;
|
||||
import org.apache.rocketmq.controller.elect.impl.DefaultElectPolicy;
|
||||
import org.apache.rocketmq.controller.impl.DLedgerController;
|
||||
import org.apache.rocketmq.controller.impl.DefaultBrokerHeartbeatManager;
|
||||
import org.apache.rocketmq.controller.processor.ControllerRequestProcessor;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.Configuration;
|
||||
import org.apache.rocketmq.remoting.RemotingClient;
|
||||
import org.apache.rocketmq.remoting.RemotingServer;
|
||||
import org.apache.rocketmq.remoting.netty.NettyClientConfig;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRemotingClient;
|
||||
import org.apache.rocketmq.remoting.netty.NettyServerConfig;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.remoting.protocol.RequestCode;
|
||||
import org.apache.rocketmq.remoting.protocol.body.BrokerMemberGroup;
|
||||
import org.apache.rocketmq.remoting.protocol.header.NotifyBrokerRoleChangedRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterResponseHeader;
|
||||
|
||||
public class ControllerManager {
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.CONTROLLER_LOGGER_NAME);
|
||||
|
||||
@@ -24,7 +24,6 @@ import io.openmessaging.storage.dledger.MemberState;
|
||||
import io.openmessaging.storage.dledger.protocol.AppendEntryRequest;
|
||||
import io.openmessaging.storage.dledger.protocol.AppendEntryResponse;
|
||||
import io.openmessaging.storage.dledger.protocol.BatchAppendEntryRequest;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -41,14 +40,6 @@ import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ServiceThread;
|
||||
import org.apache.rocketmq.common.ThreadFactoryImpl;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.CleanControllerBrokerDataRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetMetaDataResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.controller.Controller;
|
||||
import org.apache.rocketmq.controller.elect.ElectPolicy;
|
||||
import org.apache.rocketmq.controller.elect.impl.DefaultElectPolicy;
|
||||
@@ -64,6 +55,14 @@ import org.apache.rocketmq.remoting.RemotingServer;
|
||||
import org.apache.rocketmq.remoting.netty.NettyClientConfig;
|
||||
import org.apache.rocketmq.remoting.netty.NettyServerConfig;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.remoting.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.remoting.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.CleanControllerBrokerDataRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetMetaDataResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
|
||||
/**
|
||||
* The implementation of controller, based on dledger (raft).
|
||||
|
||||
+2
-4
@@ -17,7 +17,6 @@
|
||||
package org.apache.rocketmq.controller.impl;
|
||||
|
||||
import io.netty.channel.Channel;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Iterator;
|
||||
import java.util.List;
|
||||
@@ -28,7 +27,6 @@ import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import org.apache.rocketmq.common.BrokerAddrInfo;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.ThreadFactoryImpl;
|
||||
@@ -37,7 +35,7 @@ import org.apache.rocketmq.controller.BrokerHeartbeatManager;
|
||||
import org.apache.rocketmq.controller.BrokerLiveInfo;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.common.RemotingUtil;
|
||||
import org.apache.rocketmq.remoting.common.RemotingHelper;
|
||||
|
||||
public class DefaultBrokerHeartbeatManager implements BrokerHeartbeatManager {
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.CONTROLLER_LOGGER_NAME);
|
||||
@@ -78,7 +76,7 @@ public class DefaultBrokerHeartbeatManager implements BrokerHeartbeatManager {
|
||||
final Channel channel = next.getValue().getChannel();
|
||||
iterator.remove();
|
||||
if (channel != null) {
|
||||
RemotingUtil.closeChannel(channel);
|
||||
RemotingHelper.closeChannel(channel);
|
||||
}
|
||||
this.executor.submit(() ->
|
||||
notifyBrokerInActive(next.getKey().getClusterName(), next.getValue().getBrokerName(), next.getKey().getBrokerAddr(), next.getValue().getBrokerId()));
|
||||
|
||||
+1
-1
@@ -18,7 +18,7 @@ package org.apache.rocketmq.controller.impl.event;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.remoting.protocol.ResponseCode;
|
||||
|
||||
public class ControllerResult<T> {
|
||||
private final List<EventMessage> events;
|
||||
|
||||
+13
-13
@@ -29,19 +29,6 @@ import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.body.BrokerMemberGroup;
|
||||
import org.apache.rocketmq.common.protocol.body.InSyncStateData;
|
||||
import org.apache.rocketmq.common.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.CleanControllerBrokerDataRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerResponseHeader;
|
||||
import org.apache.rocketmq.controller.elect.ElectPolicy;
|
||||
import org.apache.rocketmq.controller.impl.event.AlterSyncStateSetEvent;
|
||||
import org.apache.rocketmq.controller.impl.event.ApplyBrokerIdEvent;
|
||||
@@ -52,6 +39,19 @@ import org.apache.rocketmq.controller.impl.event.EventMessage;
|
||||
import org.apache.rocketmq.controller.impl.event.EventType;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
import org.apache.rocketmq.logging.InternalLoggerFactory;
|
||||
import org.apache.rocketmq.remoting.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.remoting.protocol.body.BrokerMemberGroup;
|
||||
import org.apache.rocketmq.remoting.protocol.body.InSyncStateData;
|
||||
import org.apache.rocketmq.remoting.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.AlterSyncStateSetResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.CleanControllerBrokerDataRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerResponseHeader;
|
||||
|
||||
/**
|
||||
* The manager that manages the replicas info for all brokers. We can think of this class as the controller's memory
|
||||
|
||||
+20
-20
@@ -25,16 +25,6 @@ import java.util.concurrent.TimeUnit;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.BrokerHeartbeatRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.CleanControllerBrokerDataRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerResponseHeader;
|
||||
import org.apache.rocketmq.controller.BrokerHeartbeatManager;
|
||||
import org.apache.rocketmq.controller.ControllerManager;
|
||||
import org.apache.rocketmq.logging.InternalLogger;
|
||||
@@ -43,17 +33,27 @@ import org.apache.rocketmq.remoting.common.RemotingHelper;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRequestProcessor;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingSerializable;
|
||||
import org.apache.rocketmq.remoting.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.remoting.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.BrokerHeartbeatRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.CleanControllerBrokerDataRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerResponseHeader;
|
||||
|
||||
import static org.apache.rocketmq.common.protocol.RequestCode.BROKER_HEARTBEAT;
|
||||
import static org.apache.rocketmq.common.protocol.RequestCode.CLEAN_BROKER_DATA;
|
||||
import static org.apache.rocketmq.common.protocol.RequestCode.CONTROLLER_ALTER_SYNC_STATE_SET;
|
||||
import static org.apache.rocketmq.common.protocol.RequestCode.CONTROLLER_ELECT_MASTER;
|
||||
import static org.apache.rocketmq.common.protocol.RequestCode.CONTROLLER_GET_METADATA_INFO;
|
||||
import static org.apache.rocketmq.common.protocol.RequestCode.CONTROLLER_GET_REPLICA_INFO;
|
||||
import static org.apache.rocketmq.common.protocol.RequestCode.CONTROLLER_GET_SYNC_STATE_DATA;
|
||||
import static org.apache.rocketmq.common.protocol.RequestCode.CONTROLLER_REGISTER_BROKER;
|
||||
import static org.apache.rocketmq.common.protocol.RequestCode.GET_CONTROLLER_CONFIG;
|
||||
import static org.apache.rocketmq.common.protocol.RequestCode.UPDATE_CONTROLLER_CONFIG;
|
||||
import static org.apache.rocketmq.remoting.protocol.RequestCode.BROKER_HEARTBEAT;
|
||||
import static org.apache.rocketmq.remoting.protocol.RequestCode.CLEAN_BROKER_DATA;
|
||||
import static org.apache.rocketmq.remoting.protocol.RequestCode.CONTROLLER_ALTER_SYNC_STATE_SET;
|
||||
import static org.apache.rocketmq.remoting.protocol.RequestCode.CONTROLLER_ELECT_MASTER;
|
||||
import static org.apache.rocketmq.remoting.protocol.RequestCode.CONTROLLER_GET_METADATA_INFO;
|
||||
import static org.apache.rocketmq.remoting.protocol.RequestCode.CONTROLLER_GET_REPLICA_INFO;
|
||||
import static org.apache.rocketmq.remoting.protocol.RequestCode.CONTROLLER_GET_SYNC_STATE_DATA;
|
||||
import static org.apache.rocketmq.remoting.protocol.RequestCode.CONTROLLER_REGISTER_BROKER;
|
||||
import static org.apache.rocketmq.remoting.protocol.RequestCode.GET_CONTROLLER_CONFIG;
|
||||
import static org.apache.rocketmq.remoting.protocol.RequestCode.UPDATE_CONTROLLER_CONFIG;
|
||||
|
||||
/**
|
||||
* Processor for controller request
|
||||
|
||||
+7
-7
@@ -28,12 +28,6 @@ import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.client.exception.MQBrokerException;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.MixAll;
|
||||
import org.apache.rocketmq.common.protocol.RequestCode;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.BrokerHeartbeatRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerResponseHeader;
|
||||
import org.apache.rocketmq.controller.ControllerManager;
|
||||
import org.apache.rocketmq.controller.impl.DLedgerController;
|
||||
import org.apache.rocketmq.remoting.RemotingClient;
|
||||
@@ -41,12 +35,18 @@ import org.apache.rocketmq.remoting.netty.NettyClientConfig;
|
||||
import org.apache.rocketmq.remoting.netty.NettyRemotingClient;
|
||||
import org.apache.rocketmq.remoting.netty.NettyServerConfig;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.remoting.protocol.RequestCode;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.BrokerHeartbeatRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerResponseHeader;
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
import static org.apache.rocketmq.common.protocol.ResponseCode.CONTROLLER_NOT_LEADER;
|
||||
import static org.apache.rocketmq.remoting.protocol.RemotingSysResponseCode.SUCCESS;
|
||||
import static org.apache.rocketmq.remoting.protocol.ResponseCode.CONTROLLER_NOT_LEADER;
|
||||
import static org.awaitility.Awaitility.await;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
+10
-10
@@ -28,20 +28,20 @@ import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerResponseHeader;
|
||||
import org.apache.rocketmq.controller.Controller;
|
||||
import org.apache.rocketmq.controller.elect.impl.DefaultElectPolicy;
|
||||
import org.apache.rocketmq.controller.impl.DLedgerController;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingSerializable;
|
||||
import org.apache.rocketmq.remoting.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.remoting.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerResponseHeader;
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
@@ -314,4 +314,4 @@ public class DLedgerControllerTest {
|
||||
syncStateSet.add("127.0.0.1:9002");
|
||||
assertEquals(syncStateSetResult.getSyncStateSet(), syncStateSet);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+12
-13
@@ -19,27 +19,26 @@ package org.apache.rocketmq.controller.impl.controller.impl.manager;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
|
||||
import org.apache.rocketmq.common.ControllerConfig;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.AlterSyncStateSetResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.CleanControllerBrokerDataRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.RegisterBrokerToControllerResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.ElectMasterResponseHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.header.namesrv.controller.GetReplicaInfoResponseHeader;
|
||||
import org.apache.rocketmq.controller.elect.ElectPolicy;
|
||||
import org.apache.rocketmq.controller.elect.impl.DefaultElectPolicy;
|
||||
import org.apache.rocketmq.controller.impl.DefaultBrokerHeartbeatManager;
|
||||
import org.apache.rocketmq.controller.impl.manager.ReplicasInfoManager;
|
||||
import org.apache.rocketmq.controller.impl.event.ControllerResult;
|
||||
import org.apache.rocketmq.controller.impl.event.ElectMasterEvent;
|
||||
import org.apache.rocketmq.controller.impl.event.EventMessage;
|
||||
import org.apache.rocketmq.controller.impl.manager.ReplicasInfoManager;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingSerializable;
|
||||
import org.apache.rocketmq.remoting.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.remoting.protocol.body.SyncStateSet;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.AlterSyncStateSetRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.AlterSyncStateSetResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.CleanControllerBrokerDataRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.ElectMasterResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.GetReplicaInfoResponseHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerRequestHeader;
|
||||
import org.apache.rocketmq.remoting.protocol.header.namesrv.controller.RegisterBrokerToControllerResponseHeader;
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
Reference in New Issue
Block a user