From 84ba11929e761596186126d0fb231af611bc613b Mon Sep 17 00:00:00 2001 From: calvix <7136358+calvix@users.noreply.github.com> Date: Mon, 5 Oct 2026 08:55:08 +0200 Subject: [PATCH 1/3] jobs: count VmWork jobs when draining a management server for maintenance countPendingNonPseudoJobs() filtered on instance_type != 'Thread'. VmWork jobs (VmWorkStart, VmWorkStop, VmWorkAttachVolume, ...) carry a NULL instance_type, and NULL != 'Thread' is not true in SQL, so they were never counted. The maintenance drain then reported 0 pending jobs and moved the agents away while VM work was still running on the management server, which failed that work. Count NULL instance types as well, and use the same ownership rule as the jobs a restarting management server fails (getResetJobs): executing on this server, or queued on it and not picked up yet. --- .../framework/jobs/dao/AsyncJobDaoImpl.java | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 deletions(-) diff --git a/framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/dao/AsyncJobDaoImpl.java b/framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/dao/AsyncJobDaoImpl.java index d8385e9aecd1..df137a55c704 100644 --- a/framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/dao/AsyncJobDaoImpl.java +++ b/framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/dao/AsyncJobDaoImpl.java @@ -113,9 +113,20 @@ public AsyncJobDaoImpl() { pendingNonPseudoAsyncJobsSearch = createSearchBuilder(Long.class); pendingNonPseudoAsyncJobsSearch.select(null, SearchCriteria.Func.COUNT, pendingNonPseudoAsyncJobsSearch.entity().getId()); - pendingNonPseudoAsyncJobsSearch.and("instanceTypeNEQ", pendingNonPseudoAsyncJobsSearch.entity().getInstanceType(), SearchCriteria.Op.NEQ); + // VmWork jobs carry a NULL instance_type; "NEQ 'Thread'" alone would skip them (NULL != x is not true in SQL), + // and the maintenance drain would then finish while VM work still runs on this management server. + pendingNonPseudoAsyncJobsSearch.and().op("instanceTypeNULL", pendingNonPseudoAsyncJobsSearch.entity().getInstanceType(), SearchCriteria.Op.NULL); + pendingNonPseudoAsyncJobsSearch.or("instanceTypeNEQ", pendingNonPseudoAsyncJobsSearch.entity().getInstanceType(), SearchCriteria.Op.NEQ); + pendingNonPseudoAsyncJobsSearch.cp(); pendingNonPseudoAsyncJobsSearch.and("jobStatusEQ", pendingNonPseudoAsyncJobsSearch.entity().getStatus(), SearchCriteria.Op.EQ); - pendingNonPseudoAsyncJobsSearch.and("executingMsidIN", pendingNonPseudoAsyncJobsSearch.entity().getExecutingMsid(), SearchCriteria.Op.IN); + // Same ownership rule as the jobs a restarting management server fails (getResetJobs): executing here, or + // queued here and not picked up yet. + pendingNonPseudoAsyncJobsSearch.and().op("executingMsidIN", pendingNonPseudoAsyncJobsSearch.entity().getExecutingMsid(), SearchCriteria.Op.IN); + pendingNonPseudoAsyncJobsSearch.or().op("executingMsidNULL", pendingNonPseudoAsyncJobsSearch.entity().getExecutingMsid(), SearchCriteria.Op.NULL); + pendingNonPseudoAsyncJobsSearch.and("initMsidIN", pendingNonPseudoAsyncJobsSearch.entity().getInitMsid(), SearchCriteria.Op.IN); + pendingNonPseudoAsyncJobsSearch.cp(); + pendingNonPseudoAsyncJobsSearch.cp(); + pendingNonPseudoAsyncJobsSearch.done(); } @Override @@ -283,6 +294,7 @@ public long countPendingNonPseudoJobs(Long... msIds) { sc.setParameters("jobStatusEQ", JobInfo.Status.IN_PROGRESS); if (msIds != null) { sc.setParameters("executingMsidIN", (Object[])msIds); + sc.setParameters("initMsidIN", (Object[])msIds); } List results = customSearch(sc, null); return results.get(0); From 579141d739b363103b86525746198b096ac97329 Mon Sep 17 00:00:00 2001 From: calvix <7136358+calvix@users.noreply.github.com> Date: Mon, 5 Oct 2026 08:55:18 +0200 Subject: [PATCH 2/3] agent: move indirect agents to another management server without losing commands Preparing a management server for maintenance migrates its indirect agents with MigrateAgentConnectionCommand. The agent reconnected 3 s later whatever was still running, and the old management server treated the closed link as a failure: commands in flight lost their answers, the host went Alert with an agent investigation and HA work, and commands sent meanwhile failed with AgentUnavailableException. Under load every agent move failed operations, and once a live migration completed in libvirt while CloudStack recorded it as failed. Make the move a planned handoff: - Before the move command, put the host into Rebalancing. A disconnect while Rebalancing stays Rebalancing (no Alert, no investigation), and the agent's ping does not ask for a new startup that would pull it back. - Wait until this management server has no command queued or awaiting an answer for the host (indirect.agent.migration.idle.wait, default 600 s); otherwise leave the agent where it is. - Send the move command with send(): easySend() refuses hosts that are not Up/Connecting. - Hold new commands for a host that is Rebalancing or Connecting until it is Up on its new owner (agent.handoff.wait, default 180 s), instead of failing them. The move command itself is never held. - On the agent, finish running and queued commands before reconnecting, so their answers still go back over the current link. - Move one host at a time (indirect.agent.migration.parallelism, default 1), because the planner skips Rebalancing hosts. - If the agent does not come Up elsewhere, reconnect it instead of leaving the host Rebalancing. --- .../src/main/java/com/cloud/agent/Agent.java | 32 +++++++ .../java/com/cloud/agent/AgentManager.java | 9 ++ .../cloud/agent/manager/AgentManagerImpl.java | 57 ++++++++++- .../test/DirectAgentManagerSimpleImpl.java | 5 + .../agent/lb/IndirectAgentLBServiceImpl.java | 95 ++++++++++++++++++- 5 files changed, 193 insertions(+), 5 deletions(-) diff --git a/agent/src/main/java/com/cloud/agent/Agent.java b/agent/src/main/java/com/cloud/agent/Agent.java index 275fd41edc34..d9b9107ba20b 100644 --- a/agent/src/main/java/com/cloud/agent/Agent.java +++ b/agent/src/main/java/com/cloud/agent/Agent.java @@ -969,6 +969,7 @@ private Answer migrateAgentToOtherMS(final MigrateAgentConnectionCommand cmd) { } ScheduledExecutorService migrateAgentConnectionService = Executors.newSingleThreadScheduledExecutor(new NamedThreadFactory("MigrateAgentConnection-Job")); migrateAgentConnectionService.schedule(() -> { + waitForCommandsToFinish(MIGRATE_AGENT_CONNECTION_MAX_WAIT_SECS); migrateAgentConnection(cmd.getAvoidMsList()); }, 3, TimeUnit.SECONDS); migrateAgentConnectionService.shutdown(); @@ -980,6 +981,37 @@ private Answer migrateAgentToOtherMS(final MigrateAgentConnectionCommand cmd) { return new MigrateAgentConnectionAnswer(true); } + private static final int MIGRATE_AGENT_CONNECTION_MAX_WAIT_SECS = 300; + + /** + * Answers travel back over the current link, so a command still running or queued when the agent reconnects to + * another management server would lose its answer and fail there. Wait until nothing has run or been queued for two + * consecutive checks; the management server holds new commands for this host while it is Rebalancing. + */ + private void waitForCommandsToFinish(final int maxWaitSecs) { + final long deadline = System.currentTimeMillis() + maxWaitSecs * 1000L; + int idleChecks = 0; + while (idleChecks < 2) { + final int queued = requestHandler instanceof ThreadPoolExecutor ? ((ThreadPoolExecutor)requestHandler).getQueue().size() : 0; + idleChecks = (commandsInProgress.get() == 0 && queued == 0) ? idleChecks + 1 : 0; + if (idleChecks >= 2) { + break; + } + if (System.currentTimeMillis() >= deadline) { + logger.warn("Moving the agent connection with {} commands still in progress and {} queued after {} s", + commandsInProgress.get(), queued, maxWaitSecs); + return; + } + try { + Thread.sleep(1000); + } catch (final InterruptedException e) { + Thread.currentThread().interrupt(); + return; + } + } + logger.info("No commands in progress, moving the agent connection"); + } + private void migrateAgentConnection(List avoidMsList) { final String[] msHosts = shell.getHosts(); if (msHosts == null || msHosts.length < 1) { diff --git a/engine/components-api/src/main/java/com/cloud/agent/AgentManager.java b/engine/components-api/src/main/java/com/cloud/agent/AgentManager.java index 4d63fae33560..d118588a760c 100644 --- a/engine/components-api/src/main/java/com/cloud/agent/AgentManager.java +++ b/engine/components-api/src/main/java/com/cloud/agent/AgentManager.java @@ -54,6 +54,10 @@ public interface AgentManager { "This timeout overrides the wait global config. This holds a comma separated key value pairs containing timeout (in seconds) for specific commands. " + "For example: DhcpEntryCommand=600, SavePasswordCommand=300, VmDataCommand=300", false); + ConfigKey AgentHandoffWait = new ConfigKey<>("Advanced", Integer.class, "agent.handoff.wait", "180", + "Seconds a command to a host waits while the host's agent is moved to another management server (host in Rebalancing), " + + "instead of failing with the agent unavailable. After that the command is sent as before.", true); + ConfigKey KVMHostDiscoverySshPort = new ConfigKey<>(ConfigKey.CATEGORY_ADVANCED, Integer.class, "kvm.host.discovery.ssh.port", String.valueOf(Host.DEFAULT_SSH_PORT), "SSH port used for KVM host discovery and any other operations on host (using SSH)." + " Please note that this is applicable when port is not defined through host url while adding the KVM host.", true, ConfigKey.Scope.Cluster); @@ -155,6 +159,11 @@ enum TapAgentsAction { boolean isAgentAttached(long hostId); + /** + * @return true when no command sent from this management server to the host is queued or waiting for its answer + */ + boolean isAgentIdle(long hostId); + void disconnectWithoutInvestigation(long hostId, Status.Event event); void disconnectWithInvestigation(long hostId, Status.Event event); diff --git a/engine/orchestration/src/main/java/com/cloud/agent/manager/AgentManagerImpl.java b/engine/orchestration/src/main/java/com/cloud/agent/manager/AgentManagerImpl.java index 1215829d92f8..c30d4e8bfa31 100644 --- a/engine/orchestration/src/main/java/com/cloud/agent/manager/AgentManagerImpl.java +++ b/engine/orchestration/src/main/java/com/cloud/agent/manager/AgentManagerImpl.java @@ -75,6 +75,7 @@ import com.cloud.agent.api.Answer; import com.cloud.agent.api.CheckHealthCommand; import com.cloud.agent.api.Command; +import com.cloud.agent.api.MigrateAgentConnectionCommand; import com.cloud.agent.api.PingAnswer; import com.cloud.agent.api.PingCommand; import com.cloud.agent.api.PingRoutingCommand; @@ -654,6 +655,8 @@ public Answer[] send(final Long hostId, final Commands commands, int timeout) th final Command[] cmds = checkForCommandsAndTag(commands); + waitWhileAgentHandoff(hostId, cmds); + //check what agent is returned. final AgentAttache agent = getAttache(hostId); if (agent == null || agent.isClosed()) { @@ -703,8 +706,57 @@ protected AgentAttache getAttache(final Long hostId) throws AgentUnavailableExce return agent; } + /** + * A planned agent move (management server maintenance) puts the host into Rebalancing while its agent finishes the + * commands it is running and reconnects to another management server. Hold new commands for that host until the + * move is over, so they go to the new owner instead of failing with the agent unavailable. The command that moves + * the agent is never held. + */ + protected void waitWhileAgentHandoff(final Long hostId, final Command[] cmds) { + if (hostId == null || cmds == null) { + return; + } + for (final Command cmd : cmds) { + if (cmd instanceof MigrateAgentConnectionCommand) { + return; + } + } + HostVO host = _hostDao.findById(hostId); + if (host == null || host.getStatus() != Status.Rebalancing) { + return; + } + final long start = System.currentTimeMillis(); + final long deadline = start + AgentHandoffWait.value() * 1000L; + logger.info("Holding {} for {} until its agent has moved to another management server", cmds[0].getClass().getSimpleName(), host); + while (System.currentTimeMillis() < deadline) { + try { + Thread.sleep(500); + } catch (final InterruptedException e) { + Thread.currentThread().interrupt(); + return; + } + host = _hostDao.findById(hostId); + // Connecting is still part of the move: until the host is Up on its new owner, getAttache() here would not + // switch to a forwarding attache and the command would fail with "not in the right state: Connecting". + if (host == null || (host.getStatus() != Status.Rebalancing && host.getStatus() != Status.Connecting)) { + logger.info("Released {} for {} after {} ms, host is {}", cmds[0].getClass().getSimpleName(), host, + System.currentTimeMillis() - start, host == null ? null : host.getStatus()); + return; + } + } + logger.warn("{} is still Rebalancing after {} s ({}), sending {} anyway", host, AgentHandoffWait.value(), + AgentHandoffWait.key(), cmds[0].getClass().getSimpleName()); + } + + @Override + public boolean isAgentIdle(final long hostId) { + final AgentAttache attache = findAttache(hostId); + return attache == null || (attache.getQueueSize() == 0 && attache.getNonRecurringListenersSize() == 0); + } + @Override public long send(final Long hostId, final Commands commands, final Listener listener) throws AgentUnavailableException { + waitWhileAgentHandoff(hostId, commands.toCommands()); final AgentAttache agent = getAttache(hostId); if (agent.isClosed()) { throw new AgentUnavailableException(String.format( @@ -1713,7 +1765,8 @@ protected void processRequest(final Link link, final Request request) { logger.debug("Not processing {} for agent id={}; can't find the host in the DB", PingRoutingCommand.class.getSimpleName(), cmdHostId); } } - if (host != null && host.getStatus() != Status.Up && gatewayAccessible) { + // A Rebalancing host is moving to another management server; asking for its startup here would pull it back. + if (host != null && host.getStatus() != Status.Up && host.getStatus() != Status.Rebalancing && gatewayAccessible) { requestStartupCommand = true; } final List avoidMsList = _mshostDao.listNonUpStateMsIPs(); @@ -2134,7 +2187,7 @@ public ConfigKey[] getConfigKeys() { return new ConfigKey[] { CheckTxnBeforeSending, Workers, Port, Wait, AlertWait, DirectAgentLoadSize, DirectAgentPoolSize, DirectAgentThreadCap, EnableKVMAutoEnableDisable, ReadyCommandWait, GranularWaitTimeForCommands, RemoteAgentSslHandshakeTimeout, RemoteAgentMaxConcurrentNewConnections, - RemoteAgentNewConnectionsMonitorInterval, KVMHostDiscoverySshPort }; + RemoteAgentNewConnectionsMonitorInterval, KVMHostDiscoverySshPort, AgentHandoffWait }; } protected class SetHostParamsListener implements Listener { diff --git a/engine/storage/integration-test/src/test/java/org/apache/cloudstack/storage/test/DirectAgentManagerSimpleImpl.java b/engine/storage/integration-test/src/test/java/org/apache/cloudstack/storage/test/DirectAgentManagerSimpleImpl.java index 1d072985a667..4938df176950 100644 --- a/engine/storage/integration-test/src/test/java/org/apache/cloudstack/storage/test/DirectAgentManagerSimpleImpl.java +++ b/engine/storage/integration-test/src/test/java/org/apache/cloudstack/storage/test/DirectAgentManagerSimpleImpl.java @@ -272,6 +272,11 @@ public boolean isAgentAttached(long hostId) { return false; } + @Override + public boolean isAgentIdle(long hostId) { + return true; + } + @Override public boolean handleDirectConnectAgent(Host host, StartupCommand[] cmds, ServerResource resource, boolean forRebalance, boolean newHost) throws ConnectionException { return false; diff --git a/server/src/main/java/org/apache/cloudstack/agent/lb/IndirectAgentLBServiceImpl.java b/server/src/main/java/org/apache/cloudstack/agent/lb/IndirectAgentLBServiceImpl.java index dc7b6282b080..5320299004c3 100644 --- a/server/src/main/java/org/apache/cloudstack/agent/lb/IndirectAgentLBServiceImpl.java +++ b/server/src/main/java/org/apache/cloudstack/agent/lb/IndirectAgentLBServiceImpl.java @@ -50,8 +50,11 @@ import com.cloud.dc.DataCenterVO; import com.cloud.dc.dao.ClusterDao; import com.cloud.dc.dao.DataCenterDao; +import com.cloud.exception.AgentUnavailableException; +import com.cloud.exception.OperationTimedoutException; import com.cloud.host.Host; import com.cloud.host.HostVO; +import com.cloud.host.Status; import com.cloud.host.dao.HostDao; import com.cloud.hypervisor.Hypervisor; import com.cloud.resource.ResourceState; @@ -73,6 +76,18 @@ public class IndirectAgentLBServiceImpl extends ComponentLifecycleBase implement " Set 0 to disable it.", true, ConfigKey.Scope.Cluster); + public static final ConfigKey IndirectAgentMigrationParallelism = new ConfigKey<>("Advanced", Integer.class, + "indirect.agent.migration.parallelism", "1", + "How many indirect agents a management server going into maintenance moves to other management servers at the same time. " + + "A host is Rebalancing while its agent moves and the deployment planner skips it, so keep this low.", + true); + + public static final ConfigKey IndirectAgentMigrationIdleWait = new ConfigKey<>("Advanced", Integer.class, + "indirect.agent.migration.idle.wait", "600", + "Seconds to wait, before an indirect agent is moved to another management server, for the commands this management " + + "server already sent to it to finish. When they do not finish in time, the agent is not moved and maintenance is not prepared.", + true); + private static Map algorithmMap = new HashMap<>(); @Inject @@ -478,7 +493,7 @@ protected boolean migrateRoutingHostAgentsInCluster(long clusterId, String fromM } logger.debug(String.format("Migrating %d indirect routing host agents from management server node %d (id: %s) of zone %s, " + "cluster ID: %d", agentBasedHostsOfMsInDcAndCluster.size(), fromMsId, fromMsUuid, dc, clusterId)); - ExecutorService migrateAgentsExecutorService = Executors.newFixedThreadPool(10, new NamedThreadFactory("MigrateRoutingHostAgent-Worker")); + ExecutorService migrateAgentsExecutorService = Executors.newFixedThreadPool(Math.max(1, IndirectAgentMigrationParallelism.value()), new NamedThreadFactory("MigrateRoutingHostAgent-Worker")); Long lbCheckInterval = getLBPreferredHostCheckInterval(clusterId); boolean stopMigration = false; for (final Long hostId : agentBasedHostsOfMsInDcAndCluster) { @@ -595,21 +610,93 @@ protected void runInContext() { msList = getManagementServerList(hostId, dcId, orderedHostIdList, lbAlgorithm); } + // Planned move: Rebalancing holds new commands for this host (AgentManagerImpl.waitWhileAgentHandoff) and + // keeps the disconnect from turning the host Alert and starting an HA investigation. + final HostVO host = hostDao.findById(hostId); + final boolean handoff = host != null && host.getStatus() == Status.Up + && agentManager.agentStatusTransitTo(host, Status.Event.StartAgentRebalance, fromMsId); + if (handoff && !waitFor(() -> agentManager.isAgentIdle(hostId), IndirectAgentMigrationIdleWait.value())) { + logger.warn(String.format("Commands to host agent ID: %d did not finish in %d s, not moving its agent", hostId, IndirectAgentMigrationIdleWait.value())); + abortHandoff(hostId, fromMsId); + return; + } + final MigrateAgentConnectionCommand cmd = new MigrateAgentConnectionCommand(msList, avoidMsList, lbAlgorithm, lbCheckInterval); cmd.setWait(60); - final Answer answer = agentManager.easySend(hostId, cmd); //may not receive answer when the agent disconnects immediately and try reconnecting to other ms host + // easySend() refuses hosts that are not Up/Connecting, and the host is Rebalancing now; send() has no + // such check, and AgentManagerImpl.waitWhileAgentHandoff() never holds this command. + final Answer answer = sendMigrateCommand(hostId, cmd); //may not receive answer when the agent disconnects immediately and try reconnecting to other ms host if (answer == null) { logger.warn(String.format("Got empty answer while initiating migration of agent connection for host agent ID: %d", hostId)); } else if (!answer.getResult()) { logger.warn(String.format("Error while initiating migration of agent connection for host agent ID: %d - %s", hostId, answer.getDetails())); + if (handoff) { + abortHandoff(hostId, fromMsId); + return; + } } updateLastManagementServer(hostId, fromMsId); + if (handoff && !waitFor(() -> isHostUpElsewhere(hostId, fromMsId), AgentManager.AgentHandoffWait.value())) { + logger.warn(String.format("Host agent ID: %d is not Up on another management server %d s after its agent was asked to move, letting it reconnect", hostId, AgentManager.AgentHandoffWait.value())); + abortHandoff(hostId, fromMsId); + } } catch (final Exception e) { logger.error(String.format("Error migrating agent connection for host %d", hostId), e); } } } + private Answer sendMigrateCommand(final long hostId, final MigrateAgentConnectionCommand cmd) { + try { + return agentManager.send(hostId, cmd); + } catch (final AgentUnavailableException | OperationTimedoutException e) { + logger.warn(String.format("Unable to send %s to host agent ID: %d: %s", cmd.getClass().getSimpleName(), hostId, e.getMessage())); + return null; + } + } + + private boolean waitFor(final java.util.function.BooleanSupplier condition, final long timeoutSecs) { + final long deadline = System.currentTimeMillis() + timeoutSecs * 1000L; + while (!condition.getAsBoolean()) { + if (System.currentTimeMillis() >= deadline) { + return false; + } + try { + Thread.sleep(500); + } catch (final InterruptedException e) { + Thread.currentThread().interrupt(); + return false; + } + } + return true; + } + + private boolean isHostUpElsewhere(final long hostId, final long fromMsId) { + final HostVO host = hostDao.findById(hostId); + return host == null || (host.getStatus() == Status.Up && host.getManagementServerId() != null && host.getManagementServerId() != fromMsId); + } + + /** + * The agent did not move. There is no Rebalancing -> Up transition, so let it reconnect: reconnect() accepts a + * Rebalancing host (ShutdownRequested -> Disconnected, then Connecting -> Up), without an investigation. Only if that + * fails, mark the move failed (Disconnected); the agent's next ping brings the host back Up. + */ + private void abortHandoff(final long hostId, final long fromMsId) { + final HostVO host = hostDao.findById(hostId); + if (host == null || host.getStatus() != Status.Rebalancing) { + return; + } + try { + agentManager.reconnect(hostId); + } catch (final AgentUnavailableException | RuntimeException e) { + logger.warn(String.format("Unable to reconnect host agent ID: %d after an aborted move: %s", hostId, e.getMessage())); + final HostVO current = hostDao.findById(hostId); + if (current != null && current.getStatus() == Status.Rebalancing) { + agentManager.agentStatusTransitTo(current, Status.Event.RebalanceFailed, fromMsId); + } + } + } + private void updateLastManagementServer(long hostId, long msId) { HostVO hostVO = hostDao.findById(hostId); if (hostVO != null) { @@ -645,7 +732,9 @@ public String getConfigComponentName() { public ConfigKey[] getConfigKeys() { return new ConfigKey[] { IndirectAgentLBAlgorithm, - IndirectAgentLBCheckInterval + IndirectAgentLBCheckInterval, + IndirectAgentMigrationParallelism, + IndirectAgentMigrationIdleWait }; } } From 55e13dc7ac3c2b76fcf01711c38c849446559499 Mon Sep 17 00:00:00 2001 From: calvix <7136358+calvix@users.noreply.github.com> Date: Mon, 5 Oct 2026 08:55:18 +0200 Subject: [PATCH 3/3] cluster: reopen a peer channel the peer management server has closed When a management server restarts quickly, ClusterManagerImpl only notices the new runid ("left and rejoined quickly") and raises no left/joined event, so the other server keeps its cached channel to the old process. The first write to a channel the peer has closed still succeeds locally and its data is lost; only the next write fails with "Broken pipe" and triggers a reconnect. The lost write is typically an agent answer routed back to the restarted server, so the job waiting for it hangs until its timeout. Peer channels are only ever written to. Before reusing a cached one, do a non-blocking read: -1 means the peer closed it, so open a new connection first. --- .../manager/ClusteredAgentManagerImpl.java | 39 ++++++++++++++++++- .../ClusteredAgentManagerImplTest.java | 32 +++++++++++++++ 2 files changed, 70 insertions(+), 1 deletion(-) diff --git a/engine/orchestration/src/main/java/com/cloud/agent/manager/ClusteredAgentManagerImpl.java b/engine/orchestration/src/main/java/com/cloud/agent/manager/ClusteredAgentManagerImpl.java index 38a198b73040..7d223f65faa7 100644 --- a/engine/orchestration/src/main/java/com/cloud/agent/manager/ClusteredAgentManagerImpl.java +++ b/engine/orchestration/src/main/java/com/cloud/agent/manager/ClusteredAgentManagerImpl.java @@ -493,10 +493,47 @@ public void closePeer(final String peerName) { } } + /** + * A peer that restarted quickly gets no left/joined event (ClusterManagerImpl only updates its runid), so we keep + * a channel the peer process has already closed. The first write to it still succeeds locally and its data is lost; + * only the next write fails with "Broken pipe". Peer channels are only ever written to, so a non-blocking read that + * returns -1 means the peer has closed this connection. Inbound bytes (e.g. TLS post-handshake messages) are + * discarded; nothing unwraps on this channel. + */ + protected boolean isPeerChannelClosed(final SocketChannel ch) { + if (ch == null || !ch.isOpen() || !ch.isConnected()) { + return true; + } + if (ch.isBlocking()) { + return false; + } + try { + final ByteBuffer discard = ByteBuffer.allocate(4096); + int read; + while ((read = ch.read(discard)) > 0) { + discard.clear(); + } + return read < 0; + } catch (final IOException e) { + return true; + } + } + public SocketChannel connectToPeer(final String peerName, final SocketChannel prevCh) { synchronized (_peers) { - final SocketChannel ch = _peers.get(peerName); + SocketChannel ch = _peers.get(peerName); SSLEngine sslEngine; + if (ch != null && ch != prevCh && isPeerChannelClosed(ch)) { + logger.info("Peer management server {} has closed our connection (restarted?), opening a new one", peerName); + try { + ch.close(); + } catch (final IOException e) { + logger.debug("[ignored] failed to close stale peer channel: {}", e.getLocalizedMessage()); + } + _peers.remove(peerName); + _sslEngines.remove(peerName); + ch = null; + } if (prevCh != null) { try { prevCh.close(); diff --git a/engine/orchestration/src/test/java/com/cloud/agent/manager/ClusteredAgentManagerImplTest.java b/engine/orchestration/src/test/java/com/cloud/agent/manager/ClusteredAgentManagerImplTest.java index 5e4678f62225..a220ef69be4a 100644 --- a/engine/orchestration/src/test/java/com/cloud/agent/manager/ClusteredAgentManagerImplTest.java +++ b/engine/orchestration/src/test/java/com/cloud/agent/manager/ClusteredAgentManagerImplTest.java @@ -23,6 +23,7 @@ import com.cloud.host.Status; import com.cloud.host.dao.HostDao; import com.cloud.resource.ResourceManagerImpl; +import org.junit.Assert; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -30,12 +31,16 @@ import org.mockito.Mockito; import org.mockito.junit.MockitoJUnitRunner; +import java.net.InetSocketAddress; +import java.nio.channels.ServerSocketChannel; +import java.nio.channels.SocketChannel; import java.util.ArrayList; import java.util.List; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyBoolean; import static org.mockito.ArgumentMatchers.anyLong; +import static org.mockito.Mockito.doCallRealMethod; import static org.mockito.Mockito.doReturn; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.never; @@ -147,4 +152,31 @@ public void scanDirectAgentToLoadHostWithNonForwardAttacheAndDisconnectedTest() verify(clusteredAgentManagerImpl).investigate(agentAttache); verify(clusteredAgentManagerImpl).loadDirectlyConnectedHost(hostVO, false); } + + @Test + public void isPeerChannelClosedDetectsPeerThatClosedTheConnection() throws Exception { + ClusteredAgentManagerImpl clusteredAgentManagerImpl = mock(ClusteredAgentManagerImpl.class); + doCallRealMethod().when(clusteredAgentManagerImpl).isPeerChannelClosed(any()); + try (ServerSocketChannel server = ServerSocketChannel.open()) { + server.bind(new InetSocketAddress("127.0.0.1", 0)); + SocketChannel client = SocketChannel.open(server.getLocalAddress()); + client.configureBlocking(false); + SocketChannel accepted = server.accept(); + + Assert.assertFalse("open connection must not be reported closed", clusteredAgentManagerImpl.isPeerChannelClosed(client)); + + accepted.close(); + boolean closed = false; + for (int i = 0; i < 50 && !closed; i++) { + closed = clusteredAgentManagerImpl.isPeerChannelClosed(client); + if (!closed) { + Thread.sleep(20); + } + } + Assert.assertTrue("connection closed by the peer must be detected", closed); + client.close(); + Assert.assertTrue(clusteredAgentManagerImpl.isPeerChannelClosed(client)); + Assert.assertTrue(clusteredAgentManagerImpl.isPeerChannelClosed(null)); + } + } }