Damans227 commented on code in PR #13345: URL: https://github.com/apache/cloudstack/pull/13345#discussion_r4158490107
########## agent/src/main/java/com/cloud/agent/HostConnectProcess.java: ########## @@ -0,0 +1,355 @@ +// 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 com.cloud.agent; + +import com.cloud.agent.api.AgentConnectStatusAnswer; +import com.cloud.agent.api.AgentConnectStatusCommand; +import com.cloud.agent.api.Answer; +import com.cloud.agent.api.Command; +import com.cloud.agent.api.StartupAnswer; +import com.cloud.agent.api.StartupCommand; +import com.cloud.agent.properties.AgentProperties; +import com.cloud.agent.properties.AgentPropertiesFileHandler; +import com.cloud.agent.transport.Request; +import com.cloud.exception.CloudException; +import com.cloud.exception.OperationTimedoutException; +import com.cloud.host.Status; +import com.cloud.resource.ResourceStatusUpdater; +import com.cloud.resource.ServerResource; +import com.cloud.utils.concurrency.NamedThreadFactory; +import com.cloud.utils.nio.Link; +import org.apache.cloudstack.threadcontext.ThreadContextUtil; +import org.apache.commons.lang3.ArrayUtils; +import org.apache.logging.log4j.Logger; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.ThreadContext; + +import java.io.IOException; +import java.nio.channels.ClosedChannelException; +import java.util.Optional; +import java.util.Set; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Predicate; + +public class HostConnectProcess { + private static final Logger logger = LogManager.getLogger(HostConnectProcess.class); + + public static final int DEFAULT_ASYNC_COMMAND_TIMEOUT_SEC = + AgentPropertiesFileHandler.getPropertyValue(AgentProperties.ASYNC_COMMAND_TIMEOUT_SEC); + + public static final int DEFAULT_ASYNC_STARTUP_COMMAND_TIMEOUT_SEC = + AgentPropertiesFileHandler.getPropertyValue(AgentProperties.ASYNC_STARTUP_COMMAND_TIMEOUT_SEC); + + static final long HOST_STATUS_CHECK_INITIAL_DELAY_SEC = 10; + private long hostStatusCheckDelaySec = AgentPropertiesFileHandler.getPropertyValue(AgentProperties.AGENT_HOST_STATUS_CHECK_DELAY_SEC); + private final AtomicReference<ScheduledFuture<?>> hostStatusFutureRef = new AtomicReference<>(); + private final Agent agent; + private ScheduledExecutorService hostStatusExecutor; + + public HostConnectProcess(Agent agent) { + this.agent = agent; + initExecutors(); + } + + private void initExecutors() { + stop(); + var threadFactory = new NamedThreadFactory("Agent-" + HostStatusTask.class.getSimpleName()); + hostStatusExecutor = Executors.newScheduledThreadPool(1, threadFactory); + } + + /** + * Stops the whole connect process and cancels all scheduled asynchronous tasks. + * Returns {@link Boolean#TRUE} if {@link HostConnectProcess} was waiting for {@link StartupAnswer}. + */ + public boolean stop() { + logger.debug("Stopping connect process. The process is active: {}", isInProgress()); + stopHostStatusExecutor(); + logger.debug("Stopped executor"); + Optional<? extends ScheduledFuture<?>> hostStatusOpt = Optional.ofNullable(hostStatusFutureRef.getAndSet(null)) + .filter(Predicate.not(ScheduledFuture::isCancelled)); + + hostStatusOpt.ifPresent(future -> future.cancel(true)); + logger.debug("Cancelled future"); + + return hostStatusOpt.isPresent(); + } + + private void stopHostStatusExecutor() { + if (hostStatusExecutor != null) { + hostStatusExecutor.shutdownNow(); + hostStatusExecutor = null; + } + } + + public void scheduleConnectProcess(Link link, boolean connectionTransfer) { + logger.debug("Scheduling connect process for {}", link); + initExecutors(); + + var task = new HostStatusTask(link, connectionTransfer, agent, hostStatusFutureRef); + var future = hostStatusExecutor.scheduleWithFixedDelay(ThreadContextUtil.wrapThreadContext(task), + HOST_STATUS_CHECK_INITIAL_DELAY_SEC, + hostStatusCheckDelaySec, TimeUnit.SECONDS); + hostStatusFutureRef.set(future); + } + + /** + * Returns {@link Boolean#TRUE} if {@link HostStatusTask} created and scheduled. + * That means there is already {@link Status#Connecting} process is running. + */ + public boolean isInProgress() { + return Optional.ofNullable(hostStatusFutureRef.get()) + .filter(Predicate.not(ScheduledFuture::isCancelled)).isPresent(); + } + + public void updateHostStatusCheckDelay(int newDelaySec) { + logger.info("Updating host status check delay from {} to {} seconds", hostStatusCheckDelaySec, newDelaySec); + this.hostStatusCheckDelaySec = newDelaySec; + } + + /** + * Task wait for the Host to be available to connect to submit {@link StartupCommand}. Review Comment: got it, thanks ########## agent/src/main/java/com/cloud/agent/HostConnectProcess.java: ########## @@ -0,0 +1,355 @@ +// 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 com.cloud.agent; + +import com.cloud.agent.api.AgentConnectStatusAnswer; +import com.cloud.agent.api.AgentConnectStatusCommand; +import com.cloud.agent.api.Answer; +import com.cloud.agent.api.Command; +import com.cloud.agent.api.StartupAnswer; +import com.cloud.agent.api.StartupCommand; +import com.cloud.agent.properties.AgentProperties; +import com.cloud.agent.properties.AgentPropertiesFileHandler; +import com.cloud.agent.transport.Request; +import com.cloud.exception.CloudException; +import com.cloud.exception.OperationTimedoutException; +import com.cloud.host.Status; +import com.cloud.resource.ResourceStatusUpdater; +import com.cloud.resource.ServerResource; +import com.cloud.utils.concurrency.NamedThreadFactory; +import com.cloud.utils.nio.Link; +import org.apache.cloudstack.threadcontext.ThreadContextUtil; +import org.apache.commons.lang3.ArrayUtils; +import org.apache.logging.log4j.Logger; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.ThreadContext; + +import java.io.IOException; +import java.nio.channels.ClosedChannelException; +import java.util.Optional; +import java.util.Set; +import java.util.concurrent.Executors; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.ScheduledFuture; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Predicate; + +public class HostConnectProcess { + private static final Logger logger = LogManager.getLogger(HostConnectProcess.class); + + public static final int DEFAULT_ASYNC_COMMAND_TIMEOUT_SEC = + AgentPropertiesFileHandler.getPropertyValue(AgentProperties.ASYNC_COMMAND_TIMEOUT_SEC); + + public static final int DEFAULT_ASYNC_STARTUP_COMMAND_TIMEOUT_SEC = + AgentPropertiesFileHandler.getPropertyValue(AgentProperties.ASYNC_STARTUP_COMMAND_TIMEOUT_SEC); + + static final long HOST_STATUS_CHECK_INITIAL_DELAY_SEC = 10; + private long hostStatusCheckDelaySec = AgentPropertiesFileHandler.getPropertyValue(AgentProperties.AGENT_HOST_STATUS_CHECK_DELAY_SEC); Review Comment: ok thats consistent now, thanks ########## utils/src/main/java/org/apache/cloudstack/threadcontext/ThreadContextUtil.java: ########## @@ -0,0 +1,59 @@ +// 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.cloudstack.threadcontext; + +import org.apache.logging.log4j.ThreadContext; + +import java.util.HashMap; +import java.util.Map; + +/** + * Utility class, helps to propagate {@link ThreadContext} values from parent to child threads. + * + * @author mprokopchuk + */ +public class ThreadContextUtil { + /** + * Wrap {@link Runnable} to propagate {@link ThreadContext} values. + * + * @param delegate + * @return + */ + public static Runnable wrapThreadContext(Runnable delegate) { + @SuppressWarnings("unchecked") Review Comment: ok that covers it, thanks ########## engine/orchestration/src/main/java/com/cloud/agent/manager/AgentManagerImpl.java: ########## @@ -762,40 +985,124 @@ public void notifyMonitorsOfNewlyAddedHost(long hostId) { protected AgentAttache notifyMonitorsOfConnection(final AgentAttache attache, final StartupCommand[] cmds, final boolean forRebalance) throws ConnectionException { final long hostId = attache.getId(); final HostVO host = _hostDao.findById(hostId); - for (final Pair<Integer, Listener> monitor : _hostMonitors) { - logger.debug("Sending Connect to listener: {}, for rebalance: {}", monitor.second().getClass().getSimpleName(), forRebalance); - for (StartupCommand cmd : cmds) { - try { - logger.debug("process connection to issue: {} for host: {}, forRebalance: {}", ReflectionToStringBuilderUtils.reflectOnlySelectedFields(cmd, "id", "type", "msHostList", "connectionTransferred"), hostId, forRebalance); - monitor.second().processConnect(host, cmd, forRebalance); - } catch (final ConnectionException ce) { - if (ce.isSetupError()) { - logger.warn("Monitor {} says there is an error in the connect process for {} due to {}", monitor.second().getClass().getSimpleName(), hostId, ce.getMessage()); - handleDisconnectWithoutInvestigation(attache, Event.AgentDisconnected, true, true); - throw ce; - } else { - logger.info("Monitor {} says not to continue the connect process for {} due to {}", monitor.second().getClass().getSimpleName(), hostId, ce.getMessage()); - handleDisconnectWithoutInvestigation(attache, Event.ShutdownRequested, true, true); - return attache; - } - } catch (final HypervisorVersionChangedException hvce) { - handleDisconnectWithoutInvestigation(attache, Event.ShutdownRequested, true, true); - throw new CloudRuntimeException("Unable to connect " + attache.getId(), hvce); - } catch (final Exception e) { - logger.error("Monitor {} says there is an error in the connect process for {} due to {}", monitor.second().getClass().getSimpleName(), hostId, e.getMessage(), e); - handleDisconnectWithoutInvestigation(attache, Event.AgentDisconnected, true, true); - throw new CloudRuntimeException("Unable to connect " + attache.getId(), e); - } + Optional<HostVO> hostOpt = Optional.ofNullable(host); + String hostName = hostOpt.map(HostVO::getName).orElse(null); + String hostUuid = hostOpt.map(HostVO::getUuid).orElse(null); + boolean processSuccessful; + Profiler processProfiler = new Profiler(); + processProfiler.start(); + logger.debug("Process connect for host: {} ({}) for rebalance: {}", hostName, hostUuid, forRebalance); + try { + processSuccessful = performProcessConnect(attache, cmds, forRebalance, hostName, hostUuid, host); + } finally { + processProfiler.stop(); + long processDurationSec = processProfiler.getDurationInMillis() / 1000; + logger.debug("Finished process connect for host: {} ({}) for rebalance: {} duration, s: {}", + hostName, hostUuid, forRebalance, processDurationSec); + } + if (processSuccessful) { + sendReadyCommand(attache, host, hostId); + + agentStatusTransitTo(host, Event.Ready, _nodeId); + attache.ready(); + } + return attache; + } + + /** + * @return true if commands successfully completed + */ + private boolean performProcessConnect(AgentAttache attache, StartupCommand[] cmds, boolean forRebalance, String hostName, String hostUuid, HostVO host) throws ConnectionException { + int monitorsSize = _hostMonitors.size(); + boolean processSuccessful = true; + for (int monitorIndex = 0; monitorIndex < monitorsSize && processSuccessful; monitorIndex++) { + Pair<Integer, Listener> monitor = _hostMonitors.get(monitorIndex); + Listener listener = monitor.second(); + String monitorClassName = listener.getClass().getSimpleName(); + Profiler monitorProfiler = new Profiler(); + logger.debug("Process connect for host: {} ({}) listener {} ({} of {}) for rebalance: {}", + hostName, hostUuid, monitorClassName, monitorIndex, monitorsSize, forRebalance); + try { + monitorProfiler.start(); + processSuccessful = performProcessConnect(attache, cmds, forRebalance, hostName, hostUuid, monitorClassName, listener, host); + } finally { + monitorProfiler.stop(); + long monitorDurationSec = monitorProfiler.getDurationInMillis() / 1000; + logger.debug("Finished process connect for host: {} ({}) listener: {} ({} of {}) for rebalance: {} duration, s: {}", + hostName, hostUuid, monitorClassName, monitorIndex, monitorsSize, forRebalance, monitorDurationSec); + } + } + return processSuccessful; + } + + /** + * @return true if commands successfully completed + */ + private boolean performProcessConnect(AgentAttache attache, StartupCommand[] cmds, boolean forRebalance, String hostName, String hostUuid, String monitorClassName, Listener listener, HostVO host) throws ConnectionException { + boolean processSuccessful = true; + int commandsSize = cmds.length; + for (int commandIndex = 0; commandIndex < commandsSize && processSuccessful; commandIndex++) { + StartupCommand command = cmds[commandIndex]; + String commandClassName = command.getClass().getSimpleName(); + boolean connTransferred = command.isConnectionTransferred(); + long processStart = System.currentTimeMillis(); + try { + logger.debug("Process command: {} ({} of {}) in connect for host: {} ({}) listener: {} for rebalance: {} connection transferred: {}", + commandClassName, commandIndex, commandsSize, hostName, hostUuid, monitorClassName, forRebalance, connTransferred); + listener.processConnect(host, command, forRebalance); + } catch (Exception e) { + handleConnectionException(attache, e, processStart, host, monitorClassName); + processSuccessful = false; + } finally { + long commandDurationSec = (System.currentTimeMillis() - processStart) / 1000; + logger.debug("Finished processing command: {} ({} of {}) in connect for host: {} ({}) listener: {} for rebalance: {} connection transferred: {} duration, s: {}", + commandClassName, commandIndex, commandsSize, hostName, hostUuid, monitorClassName, forRebalance, connTransferred, commandDurationSec); + } + } + return processSuccessful; + } + + private void handleConnectionException(AgentAttache attache, Exception e, long processStart, HostVO host, String monitorClassName) throws ConnectionException { + long processDurationMs = System.currentTimeMillis() - processStart; + StringBuilder summaryBuilder = getSummaryMsgBuilder("Failed to connect Host", + EventTypes.EVENT_HOST_RECONNECT, host.getUuid(), host.getName(), null, + e.getMessage(), monitorClassName, processDurationMs); + logger.fatal(summaryBuilder.toString(), e); + + // add alert to the Host here + if (e instanceof ConnectionException) { + ConnectionException ce = (ConnectionException) e; + // XXX: in case of Storage Pool issue we are ending up here Review Comment: got it, thanks ########## framework/db/src/main/java/com/cloud/utils/db/GlobalLock.java: ########## @@ -45,127 +43,248 @@ * </p> */ public class GlobalLock { - protected Logger logger = LogManager.getLogger(getClass()); + protected final static Logger logger = LogManager.getLogger(GlobalLock.class); private String name; - private int lockCount = 0; - private Thread ownerThread = null; - - private int referenceCount = 0; - private long holdingStartTick = 0; - - private static Map<String, GlobalLock> s_lockMap = new HashMap<String, GlobalLock>(); + /** + * DB lock count. + * Increments on {@link GlobalLock#lock(int)} and decrements on {@link GlobalLock#unlock()}. + * Upon {@link GlobalLock#unlock()}, if {@link GlobalLock#lockCount} is less than 1, then lock removed from DB + */ + private int lockCount; + + /** + * Internal (in-memory) lock count. + * Increments on {@link GlobalLock#addRef()} and indirectly on {@link GlobalLock#getInternLock(String)} and + * decrements on {@link GlobalLock#releaseRef()}, {@link GlobalLock#unlock()} and on {@link GlobalLock#lock(int)} + * if DB lock is unsuccessful + */ + private int referenceCount; + + /** + * Thread that owns lock. If lock called from different thread, it will be waiting for the owner to unlock it + * within requested timeout. If owner thread call {@link GlobalLock#lock(int)} again, then + * {@link GlobalLock#lockCount} will be incremented. + * If {@link GlobalLock#unlock()} called by owner thread, or DB lock will be unsuccessful, then owner thread will be + * nullified. + */ + private Thread ownerThread; + + /** + * Variable to hold lock duration in milliseconds. Used for information only. + */ + private long holdingStartTick; + + /** + * Holds all created locks. + */ + private static Map<String, GlobalLock> s_lockMap = new HashMap<>(); + + /** + * Create lock. + * + * @param name lock name + */ private GlobalLock(String name) { this.name = name; } + /** + * Increment reference count to lock. + * + * @return reference count + */ public int addRef() { synchronized (this) { referenceCount++; return referenceCount; } } + /** + * Decrement reference count to lock. + * + * @return reference count + */ public int releaseRef() { - int refCount; - boolean needToRemove = false; synchronized (this) { + if (logger.isDebugEnabled()) { + logger.debug("Releasing reference for internal lock {}, reference count: {}, lock count: {}", + name, referenceCount, lockCount); + } referenceCount--; - refCount = referenceCount; - - if (referenceCount < 0) - logger.warn("Unmatched Global lock " + name + " reference usage detected, check your code!"); - if (referenceCount == 0) + if (referenceCount < 0) { + logger.warn("Unmatched internal lock {} reference usage detected (reference count: {}, " + + "lock count: {}), check your code!", name, referenceCount, lockCount); + } else if (referenceCount < 1) { needToRemove = true; + } } - if (needToRemove) + if (needToRemove) { + if (logger.isDebugEnabled()) { + logger.debug("Need to release internal lock {}", name); + } releaseInternLock(name); + } + if (logger.isDebugEnabled()) { + logger.debug("Released reference for lock {}, reference count: {}", name, referenceCount); + } + return referenceCount; + } - return refCount; + public static boolean isLockAvailable(String name) { + if (logger.isDebugEnabled()) { + logger.debug("Checking lock availability for {}", name); + } + boolean result = false; + try { + result = DbUtil.isFreeLock(name); + } finally { + if (logger.isDebugEnabled()) { + logger.debug("Result of checking lock availability for {}: {}", name, result); + } + } + return result; } + /** + * Registers internal lock (in memory) object. Does not create any lock in DB yet. + * + * @param name lock name + * @return lock object + */ public static GlobalLock getInternLock(String name) { synchronized (s_lockMap) { + GlobalLock lock; if (s_lockMap.containsKey(name)) { - GlobalLock lock = s_lockMap.get(name); - lock.addRef(); - return lock; + lock = s_lockMap.get(name); + if (logger.isDebugEnabled()) { + logger.debug("Internal lock {} already exists with reference count {} and lock count {}", + name, lock.referenceCount, lock.lockCount); + } } else { - GlobalLock lock = new GlobalLock(name); - lock.addRef(); + lock = new GlobalLock(name); + if (logger.isDebugEnabled()) { + logger.debug("Internal lock {} does not exist, adding", name); + } s_lockMap.put(name, lock); - return lock; } + lock.addRef(); + if (logger.isDebugEnabled()) { + logger.debug("Added reference to internal lock {}, reference count {}, lock count {}", + name, lock.referenceCount, lock.lockCount); + } + return lock; } } + /** + * Unregister internal lock (in memory) object. Does not remove any lock from DB. + * + * @param name lock name + */ private void releaseInternLock(String name) { synchronized (s_lockMap) { GlobalLock lock = s_lockMap.get(name); if (lock != null) { - if (lock.referenceCount == 0) + if (lock.referenceCount == 0) { + if (logger.isDebugEnabled()) { + logger.debug("Released internal lock {}", name); + } s_lockMap.remove(name); + } else { + if (logger.isDebugEnabled()) { + logger.debug("Not releasing internal lock {} as it has references count: {}, lock count: {}", + name, lock.referenceCount, lock.lockCount); + } + } } else { - logger.warn("Releasing " + name + ", but it is already released."); + logger.warn("Internal lock {} already released", name); } } } + /** + * Acquire or join existing DB lock. + * + * @param timeoutSeconds time in seconds during which lock needs to be obtained (it is not the lock duration) + * @return true if lock successfully obtained + */ public boolean lock(int timeoutSeconds) { int remainingMilliSeconds = timeoutSeconds * 1000; Profiler profiler = new Profiler(); boolean interrupted = false; try { while (true) { synchronized (this) { - if (ownerThread != null && ownerThread == Thread.currentThread()) { - logger.warn("Global lock re-entrance detected"); - + if (ownerThread == Thread.currentThread()) { + logger.warn("Global lock {} re-entrance detected, owner thread: {}, reference count: {}, " + + "lock count: {}", name, getThreadName(ownerThread), referenceCount, lockCount); + // if it is re-entrance, then we may have more lock counts than needed? lockCount++; - if (logger.isTraceEnabled()) - logger.trace("lock " + name + " is acquired, lock count :" + lockCount); + if (logger.isDebugEnabled()) { + logger.debug("Global lock {} joined, reference count: {}, lock count: {}", + name, referenceCount, lockCount); + } return true; - } - - if (ownerThread != null) { + } else if (ownerThread != null) { profiler.start(); try { - wait((timeoutSeconds) * 1000L); + logger.debug("Waiting {} seconds to acquire global lock {}", timeoutSeconds, name); + wait(timeoutSeconds * 1000L); } catch (InterruptedException e) { interrupted = true; } profiler.stop(); remainingMilliSeconds -= profiler.getDurationInMillis(); - if (remainingMilliSeconds < 0) + if (remainingMilliSeconds < 0) { + logger.warn("Timeout of {} seconds to acquire global lock {} has been reached, " + + "owner thread {}, reference count: {}, lock count: {}", timeoutSeconds, name, getThreadName(ownerThread), referenceCount, lockCount); return false; + } continue; } else { // take ownership temporarily to prevent others enter into stage of acquiring DB lock ownerThread = Thread.currentThread(); + // XXX: do we need it here (???) Review Comment: got it, thanks ########## server/src/main/java/org/apache/cloudstack/agent/lb/IndirectAgentLBServiceImpl.java: ########## @@ -446,6 +446,7 @@ protected boolean migrateNonRoutingHostAgentsInZone(String fromMsUuid, long from break; } + // FIXME: it is fire and forget, Management Server will never know if task failed Review Comment: ok thanks -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
