refactored code / addressed comments

This commit is contained in:
Suresh Kumar Anaparti 2025-03-10 17:52:11 +05:30
parent 9ef1c129d1
commit 515e996292
No known key found for this signature in database
GPG Key ID: D7CEAE3A9E71D0AA
5 changed files with 82 additions and 73 deletions

View File

@ -1011,16 +1011,20 @@ public class Agent implements HandlerFactory, IAgentControl, AgentStatusUpdater
listener.processControlResponse(response, (AgentControlAnswer)answer);
}
} else if (answer instanceof PingAnswer) {
if ((((PingAnswer) answer).isSendStartup()) && reconnectAllowed) {
logger.info("Management server requested startup command to reinitialize the agent");
sendStartup(link);
}
shell.setAvoidHosts(((PingAnswer) answer).getAvoidMsList());
processPingAnswer((PingAnswer) answer);
} else {
updateLastPingResponseTime();
}
}
private void processPingAnswer(final PingAnswer answer) {
if ((answer.isSendStartup()) && reconnectAllowed) {
logger.info("Management server requested startup command to reinitialize the agent");
sendStartup(link);
}
shell.setAvoidHosts(answer.getAvoidMsList());
}
public void processReadyCommand(final Command cmd) {
final ReadyCommand ready = (ReadyCommand)cmd;
// Set human readable sizes;

View File

@ -255,9 +255,7 @@ public class AgentManagerImpl extends ManagerBase implements AgentManager, Handl
_executor = new ThreadPoolExecutor(agentTaskThreads, agentTaskThreads, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(), new NamedThreadFactory("AgentTaskPool"));
_connectExecutor = new ThreadPoolExecutor(100, 500, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(), new NamedThreadFactory("AgentConnectTaskPool"));
// allow core threads to time out even when there are no items in the queue
_connectExecutor.allowCoreThreadTimeOut(true);
initConnectExecutor();
maxConcurrentNewAgentConnections = RemoteAgentMaxConcurrentNewConnections.value();
@ -273,10 +271,6 @@ public class AgentManagerImpl extends ManagerBase implements AgentManager, Handl
logger.debug("Created DirectAgentAttache pool with size: {}.", directAgentPoolSize);
_directAgentThreadCap = Math.round(directAgentPoolSize * DirectAgentThreadCap.value()) + 1; // add 1 to always make the value > 0
_monitorExecutor = new ScheduledThreadPoolExecutor(1, new NamedThreadFactory("AgentMonitor"));
newAgentConnectionsMonitor = Executors.newScheduledThreadPool(1, new NamedThreadFactory("NewAgentConnectionsMonitor"));
initializeCommandTimeouts();
return true;
@ -388,10 +382,8 @@ public class AgentManagerImpl extends ManagerBase implements AgentManager, Handl
public void onManagementServerCancelMaintenance() {
logger.debug("Management server maintenance disabled");
if (_connectExecutor.isShutdown()) {
_connectExecutor = new ThreadPoolExecutor(100, 500, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(), new NamedThreadFactory("AgentConnectTaskPool"));
_connectExecutor.allowCoreThreadTimeOut(true);
initConnectExecutor();
}
startDirectlyConnectedHosts(true);
if (_connection != null) {
try {
@ -402,16 +394,30 @@ public class AgentManagerImpl extends ManagerBase implements AgentManager, Handl
}
if (_monitorExecutor.isShutdown()) {
_monitorExecutor = new ScheduledThreadPoolExecutor(1, new NamedThreadFactory("AgentMonitor"));
_monitorExecutor.scheduleWithFixedDelay(new MonitorTask(), mgmtServiceConf.getPingInterval(), mgmtServiceConf.getPingInterval(), TimeUnit.SECONDS);
initAndScheduleMonitorExecutor();
}
if (newAgentConnectionsMonitor.isShutdown()) {
final int cleanupTimeInSecs = Wait.value();
newAgentConnectionsMonitor = Executors.newScheduledThreadPool(1, new NamedThreadFactory("NewAgentConnectionsMonitor"));
newAgentConnectionsMonitor.scheduleAtFixedRate(new AgentNewConnectionsMonitorTask(), cleanupTimeInSecs, cleanupTimeInSecs, TimeUnit.SECONDS);
initAndScheduleAgentConnectionsMonitor();
}
}
private void initConnectExecutor() {
_connectExecutor = new ThreadPoolExecutor(100, 500, 60L, TimeUnit.SECONDS, new LinkedBlockingQueue<>(), new NamedThreadFactory("AgentConnectTaskPool"));
// allow core threads to time out even when there are no items in the queue
_connectExecutor.allowCoreThreadTimeOut(true);
}
private void initAndScheduleMonitorExecutor() {
_monitorExecutor = new ScheduledThreadPoolExecutor(1, new NamedThreadFactory("AgentMonitor"));
_monitorExecutor.scheduleWithFixedDelay(new MonitorTask(), mgmtServiceConf.getPingInterval(), mgmtServiceConf.getPingInterval(), TimeUnit.SECONDS);
}
private void initAndScheduleAgentConnectionsMonitor() {
final int cleanupTimeInSecs = Wait.value();
newAgentConnectionsMonitor = Executors.newScheduledThreadPool(1, new NamedThreadFactory("NewAgentConnectionsMonitor"));
newAgentConnectionsMonitor.scheduleAtFixedRate(new AgentNewConnectionsMonitorTask(), cleanupTimeInSecs, cleanupTimeInSecs, TimeUnit.SECONDS);
}
private AgentControlAnswer handleControlCommand(final AgentAttache attache, final AgentControlCommand cmd) {
AgentControlAnswer answer;
@ -805,12 +811,8 @@ public class AgentManagerImpl extends ManagerBase implements AgentManager, Handl
}
}
_monitorExecutor.scheduleWithFixedDelay(new MonitorTask(), mgmtServiceConf.getPingInterval(), mgmtServiceConf.getPingInterval(), TimeUnit.SECONDS);
final int cleanupTimeInSecs = Wait.value();
newAgentConnectionsMonitor.scheduleAtFixedRate(new AgentNewConnectionsMonitorTask(), cleanupTimeInSecs,
cleanupTimeInSecs, TimeUnit.SECONDS);
initAndScheduleMonitorExecutor();
initAndScheduleAgentConnectionsMonitor();
return true;
}

View File

@ -343,10 +343,7 @@ public class ManagementServerMaintenanceManagerImpl extends ManagerBase implemen
throw new CloudRuntimeException("Management server is not in the right state to prepare for shutdown");
}
final List<ManagementServerHostVO> preparingForMaintenanceOrShutDownMsList = msHostDao.listBy(State.PreparingForMaintenance, State.PreparingForShutDown);
if (CollectionUtils.isNotEmpty(preparingForMaintenanceOrShutDownMsList)) {
throw new CloudRuntimeException("Cannot prepare for shutdown, there are other management servers preparing for maintenance/shutdown");
}
checkAnyMsInPreparingStates("prepare for shutdown");
final Command[] cmds = new Command[1];
cmds[0] = new PrepareForShutdownManagementServerHostCommand(msHost.getMsid());
@ -370,10 +367,7 @@ public class ManagementServerMaintenanceManagerImpl extends ManagerBase implemen
}
if (State.Up.equals(msHost.getState())) {
final List<ManagementServerHostVO> preparingForMaintenanceOrShutDownMsList = msHostDao.listBy(State.PreparingForMaintenance, State.PreparingForShutDown);
if (CollectionUtils.isNotEmpty(preparingForMaintenanceOrShutDownMsList)) {
throw new CloudRuntimeException("Cannot trigger shutdown now, there are other management servers preparing for maintenance/shutdown");
}
checkAnyMsInPreparingStates("trigger shutdown");
msHostDao.updateState(msHost.getId(), State.PreparingForShutDown);
}
@ -430,10 +424,7 @@ public class ManagementServerMaintenanceManagerImpl extends ManagerBase implemen
throw new CloudRuntimeException("Management server is not in the right state to prepare for maintenance");
}
final List<ManagementServerHostVO> preparingForMaintenanceOrShutDownMsList = msHostDao.listBy(State.PreparingForMaintenance, State.PreparingForShutDown);
if (CollectionUtils.isNotEmpty(preparingForMaintenanceOrShutDownMsList)) {
throw new CloudRuntimeException("Cannot prepare for maintenance, there are other management servers preparing for maintenance/shutdown");
}
checkAnyMsInPreparingStates("prepare for maintenance");
if (indirectAgentLB.haveAgentBasedHosts(msHost.getMsid())) {
List<String> indirectAgentMsList = indirectAgentLB.getManagementServerList();
@ -505,6 +496,13 @@ public class ManagementServerMaintenanceManagerImpl extends ManagerBase implemen
msHostDao.updateState(msHost.getId(), State.Up);
}
private void checkAnyMsInPreparingStates(String operation) {
final List<ManagementServerHostVO> preparingForMaintenanceOrShutDownMsList = msHostDao.listBy(State.PreparingForMaintenance, State.PreparingForShutDown);
if (CollectionUtils.isNotEmpty(preparingForMaintenanceOrShutDownMsList)) {
throw new CloudRuntimeException(String.format("Cannot %s, there are other management servers preparing for maintenance/shutdown", operation));
}
}
private ManagementServerMaintenanceResponse prepareMaintenanceResponse(Long managementServerId) {
ManagementServerHostVO msHost;
Long[] msIds;

View File

@ -488,26 +488,32 @@ public class IndirectAgentLBServiceImpl extends ComponentLifecycleBase implement
List<DataCenterVO> dataCenterList = dcDao.listAll();
for (DataCenterVO dc : dataCenterList) {
List<Long> orderedHostIdList = getOrderedHostIdList(dc.getId());
if (!migrateNonRoutingHostAgentsInZone(fromMsUuid, fromMsId, dc, migrationStartTimeInMs,
timeoutDurationInMs, avoidMsList, lbAlgorithm, lbAlgorithmChanged, orderedHostIdList)) {
if (!migrateAgentsInZone(dc, fromMsUuid, fromMsId, avoidMsList, lbAlgorithm, lbAlgorithmChanged,
migrationStartTimeInMs, timeoutDurationInMs)) {
return false;
}
List<Long> clusterIds = clusterDao.listAllClusterIds(dc.getId());
if (CollectionUtils.isEmpty(clusterIds)) {
continue;
}
for (Long clusterId : clusterIds) {
if (!migrateRoutingHostAgentsInCluster(clusterId, fromMsUuid, fromMsId, dc, migrationStartTimeInMs,
timeoutDurationInMs, avoidMsList, lbAlgorithm, lbAlgorithmChanged, orderedHostIdList)) {
return false;
}
}
}
return true;
}
private boolean migrateAgentsInZone(DataCenterVO dc, String fromMsUuid, long fromMsId, List<String> avoidMsList,
String lbAlgorithm, boolean lbAlgorithmChanged, long migrationStartTimeInMs, long timeoutDurationInMs) {
List<Long> orderedHostIdList = getOrderedHostIdList(dc.getId());
if (!migrateNonRoutingHostAgentsInZone(fromMsUuid, fromMsId, dc, migrationStartTimeInMs,
timeoutDurationInMs, avoidMsList, lbAlgorithm, lbAlgorithmChanged, orderedHostIdList)) {
return false;
}
List<Long> clusterIds = clusterDao.listAllClusterIds(dc.getId());
for (Long clusterId : clusterIds) {
if (!migrateRoutingHostAgentsInCluster(clusterId, fromMsUuid, fromMsId, dc, migrationStartTimeInMs,
timeoutDurationInMs, avoidMsList, lbAlgorithm, lbAlgorithmChanged, orderedHostIdList)) {
return false;
}
}
return true;
}
private final class MigrateAgentConnectionTask extends ManagedContextRunnable {
private long fromMsId;
Long hostId;

View File

@ -94,17 +94,8 @@ public abstract class NioConnection implements Callable<Boolean> {
_workers = workers;
_factory = factory;
this.factoryMaxNewConnectionsCount = factory.getMaxConcurrentNewConnectionsCount();
_executor = new ThreadPoolExecutor(_workers, 5 * _workers, 1, TimeUnit.DAYS,
new LinkedBlockingQueue<>(5 * _workers), new NamedThreadFactory(_name + "-Handler"),
new ThreadPoolExecutor.AbortPolicy());
String sslHandshakeHandlerName = _name + "-SSLHandshakeHandler";
if (factoryMaxNewConnectionsCount > 0) {
_sslHandshakeExecutor = new ThreadPoolExecutor(0, this.factoryMaxNewConnectionsCount, 30,
TimeUnit.MINUTES, new SynchronousQueue<>(), new NamedThreadFactory(sslHandshakeHandlerName),
new ThreadPoolExecutor.AbortPolicy());
} else {
_sslHandshakeExecutor = Executors.newCachedThreadPool(new NamedThreadFactory(sslHandshakeHandlerName));
}
initWorkersExecutor();
initSSLHandshakeExecutor();
}
public void setCAService(final CAService caService) {
@ -129,19 +120,10 @@ public abstract class NioConnection implements Callable<Boolean> {
_isStartup = true;
if (_executor.isShutdown()) {
_executor = new ThreadPoolExecutor(_workers, 5 * _workers, 1, TimeUnit.DAYS,
new LinkedBlockingQueue<>(5 * _workers), new NamedThreadFactory(_name + "-Handler"),
new ThreadPoolExecutor.AbortPolicy());
initWorkersExecutor();
}
if (_sslHandshakeExecutor.isShutdown()) {
String sslHandshakeHandlerName = _name + "-SSLHandshakeHandler";
if (factoryMaxNewConnectionsCount > 0) {
_sslHandshakeExecutor = new ThreadPoolExecutor(0, this.factoryMaxNewConnectionsCount, 30,
TimeUnit.MINUTES, new SynchronousQueue<>(), new NamedThreadFactory(sslHandshakeHandlerName),
new ThreadPoolExecutor.AbortPolicy());
} else {
_sslHandshakeExecutor = Executors.newCachedThreadPool(new NamedThreadFactory(sslHandshakeHandlerName));
}
initSSLHandshakeExecutor();
}
_threadExecutor = Executors.newSingleThreadExecutor(new NamedThreadFactory(this._name + "-NioConnectionHandler"));
_isRunning = true;
@ -160,6 +142,23 @@ public abstract class NioConnection implements Callable<Boolean> {
}
}
private void initWorkersExecutor() {
_executor = new ThreadPoolExecutor(_workers, 5 * _workers, 1, TimeUnit.DAYS,
new LinkedBlockingQueue<>(5 * _workers), new NamedThreadFactory(_name + "-Handler"),
new ThreadPoolExecutor.AbortPolicy());
}
private void initSSLHandshakeExecutor() {
String sslHandshakeHandlerName = _name + "-SSLHandshakeHandler";
if (factoryMaxNewConnectionsCount > 0) {
_sslHandshakeExecutor = new ThreadPoolExecutor(0, this.factoryMaxNewConnectionsCount, 30,
TimeUnit.MINUTES, new SynchronousQueue<>(), new NamedThreadFactory(sslHandshakeHandlerName),
new ThreadPoolExecutor.AbortPolicy());
} else {
_sslHandshakeExecutor = Executors.newCachedThreadPool(new NamedThreadFactory(sslHandshakeHandlerName));
}
}
public boolean isRunning() {
return !_futureTask.isDone();
}