mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
Pass the CI after merge develop
This commit is contained in:
+5
-4
@@ -168,7 +168,7 @@ public class ReplicasInfoManager {
|
||||
|
||||
// Try elect a master in syncStateSet
|
||||
if (syncStateSet.size() > 1) {
|
||||
boolean electSuccess = tryElectMaster(result, brokerName, syncStateSet, (candidate) ->
|
||||
boolean electSuccess = tryElectMaster(result, brokerName, syncStateSet, candidate ->
|
||||
!candidate.equals(syncStateInfo.getMasterAddress()) && brokerAlivePredicate.test(brokerInfo.getClusterName(), candidate));
|
||||
if (electSuccess) {
|
||||
return result;
|
||||
@@ -177,7 +177,7 @@ public class ReplicasInfoManager {
|
||||
|
||||
// Try elect a master in lagging replicas if enableElectUncleanMaster = true
|
||||
if (controllerConfig.isEnableElectUncleanMaster()) {
|
||||
boolean electSuccess = tryElectMaster(result, brokerName, brokerInfo.getAllBroker(), (candidate) ->
|
||||
boolean electSuccess = tryElectMaster(result, brokerName, brokerInfo.getAllBroker(), candidate ->
|
||||
!candidate.equals(syncStateInfo.getMasterAddress()) && brokerAlivePredicate.test(brokerInfo.getClusterName(), candidate));
|
||||
if (electSuccess) {
|
||||
return result;
|
||||
@@ -227,14 +227,15 @@ public class ReplicasInfoManager {
|
||||
final BrokerMemberGroup group = new BrokerMemberGroup(brokerInfo.getClusterName(), brokerName);
|
||||
final HashMap<String, Long> brokerIdTable = brokerInfo.getBrokerIdTable();
|
||||
final HashMap<Long, String> memberGroup = new HashMap<>();
|
||||
brokerIdTable.forEach((addr, id)->memberGroup.put(id, addr));
|
||||
brokerIdTable.forEach((addr, id) -> memberGroup.put(id, addr));
|
||||
group.setBrokerAddrs(memberGroup);
|
||||
return group;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
public ControllerResult<RegisterBrokerToControllerResponseHeader> registerBroker(final RegisterBrokerToControllerRequestHeader request) {
|
||||
public ControllerResult<RegisterBrokerToControllerResponseHeader> registerBroker(
|
||||
final RegisterBrokerToControllerRequestHeader request) {
|
||||
final String brokerName = request.getBrokerName();
|
||||
final String brokerAddress = request.getBrokerAddress();
|
||||
final ControllerResult<RegisterBrokerToControllerResponseHeader> result = new ControllerResult<>(new RegisterBrokerToControllerResponseHeader());
|
||||
|
||||
+5
-5
@@ -71,7 +71,7 @@ public class ControllerRequestProcessor implements NettyRequestProcessor {
|
||||
}
|
||||
switch (request.getCode()) {
|
||||
case CONTROLLER_ALTER_SYNC_STATE_SET: {
|
||||
final AlterSyncStateSetRequestHeader controllerRequest = request.decodeCommandCustomHeader(AlterSyncStateSetRequestHeader.class);
|
||||
final AlterSyncStateSetRequestHeader controllerRequest = (AlterSyncStateSetRequestHeader) request.decodeCommandCustomHeader(AlterSyncStateSetRequestHeader.class);
|
||||
final SyncStateSet syncStateSet = RemotingSerializable.decode(request.getBody(), SyncStateSet.class);
|
||||
final CompletableFuture<RemotingCommand> future = this.controller.alterSyncStateSet(controllerRequest, syncStateSet);
|
||||
if (future != null) {
|
||||
@@ -80,7 +80,7 @@ public class ControllerRequestProcessor implements NettyRequestProcessor {
|
||||
break;
|
||||
}
|
||||
case CONTROLLER_ELECT_MASTER: {
|
||||
final ElectMasterRequestHeader controllerRequest = request.decodeCommandCustomHeader(ElectMasterRequestHeader.class);
|
||||
final ElectMasterRequestHeader controllerRequest = (ElectMasterRequestHeader) request.decodeCommandCustomHeader(ElectMasterRequestHeader.class);
|
||||
final CompletableFuture<RemotingCommand> future = this.controller.electMaster(controllerRequest);
|
||||
if (future != null) {
|
||||
return future.get(WAIT_TIMEOUT_OUT, TimeUnit.SECONDS);
|
||||
@@ -88,7 +88,7 @@ public class ControllerRequestProcessor implements NettyRequestProcessor {
|
||||
break;
|
||||
}
|
||||
case CONTROLLER_REGISTER_BROKER: {
|
||||
final RegisterBrokerToControllerRequestHeader controllerRequest = request.decodeCommandCustomHeader(RegisterBrokerToControllerRequestHeader.class);
|
||||
final RegisterBrokerToControllerRequestHeader controllerRequest = (RegisterBrokerToControllerRequestHeader) request.decodeCommandCustomHeader(RegisterBrokerToControllerRequestHeader.class);
|
||||
final CompletableFuture<RemotingCommand> future = this.controller.registerBroker(controllerRequest);
|
||||
if (future != null) {
|
||||
final RemotingCommand response = future.get(WAIT_TIMEOUT_OUT, TimeUnit.SECONDS);
|
||||
@@ -102,7 +102,7 @@ public class ControllerRequestProcessor implements NettyRequestProcessor {
|
||||
break;
|
||||
}
|
||||
case CONTROLLER_GET_REPLICA_INFO: {
|
||||
final GetReplicaInfoRequestHeader controllerRequest = request.decodeCommandCustomHeader(GetReplicaInfoRequestHeader.class);
|
||||
final GetReplicaInfoRequestHeader controllerRequest = (GetReplicaInfoRequestHeader) request.decodeCommandCustomHeader(GetReplicaInfoRequestHeader.class);
|
||||
final CompletableFuture<RemotingCommand> future = this.controller.getReplicaInfo(controllerRequest);
|
||||
if (future != null) {
|
||||
return future.get(WAIT_TIMEOUT_OUT, TimeUnit.SECONDS);
|
||||
@@ -113,7 +113,7 @@ public class ControllerRequestProcessor implements NettyRequestProcessor {
|
||||
return this.controller.getControllerMetadata();
|
||||
}
|
||||
case BROKER_HEARTBEAT: {
|
||||
final BrokerHeartbeatRequestHeader requestHeader = request.decodeCommandCustomHeader(BrokerHeartbeatRequestHeader.class);
|
||||
final BrokerHeartbeatRequestHeader requestHeader = (BrokerHeartbeatRequestHeader) request.decodeCommandCustomHeader(BrokerHeartbeatRequestHeader.class);
|
||||
this.heartbeatManager.onBrokerHeartbeat(requestHeader.getClusterName(), requestHeader.getBrokerAddr());
|
||||
return RemotingCommand.createResponseCommand(ResponseCode.SUCCESS, "Heart beat success");
|
||||
}
|
||||
|
||||
+2
-2
@@ -127,7 +127,7 @@ public class ControllerManagerTest {
|
||||
assert response != null;
|
||||
switch (response.getCode()) {
|
||||
case SUCCESS: {
|
||||
return response.decodeCommandCustomHeader(RegisterBrokerToControllerResponseHeader.class);
|
||||
return (RegisterBrokerToControllerResponseHeader) response.decodeCommandCustomHeader(RegisterBrokerToControllerResponseHeader.class);
|
||||
}
|
||||
case CONTROLLER_NOT_LEADER: {
|
||||
throw new MQBrokerException(response.getCode(), "Controller leader was changed");
|
||||
@@ -175,7 +175,7 @@ public class ControllerManagerTest {
|
||||
final GetReplicaInfoRequestHeader requestHeader = new GetReplicaInfoRequestHeader("broker1");
|
||||
final RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.CONTROLLER_GET_REPLICA_INFO, requestHeader);
|
||||
final RemotingCommand response = this.remotingClient1.invokeSync(leaderAddr, request, 3000);
|
||||
final GetReplicaInfoResponseHeader responseHeader = response.decodeCommandCustomHeader(GetReplicaInfoResponseHeader.class);
|
||||
final GetReplicaInfoResponseHeader responseHeader = (GetReplicaInfoResponseHeader) response.decodeCommandCustomHeader(GetReplicaInfoResponseHeader.class);
|
||||
assertEquals(responseHeader.getMasterAddress(), "127.0.0.1:8001");
|
||||
|
||||
executor.shutdown();
|
||||
|
||||
Reference in New Issue
Block a user