Background

VM life cycle in CloudStack is current represented through a number of lifecycle VM states, following is a complete list of these states

 Starting, Running, Stopping, Stopped, Destroyed, Migrating, Expunging, Error, Unknown

Compared with VM states defined in underlying hypervisors, CloudStack lifecycle VM states contain more information that is to reflect VM's cloud environment, when we say a CloudStack VM is running, it usually means that

To manage a CloudStack VM properly, current CloudStack has hypervisor resource-agent to participate VM lifecycle state management and periodically sync-back with CloudStack management server. Therefore, in addition for hypervisor resource agent to be aware of hypervisor specific VM power state, it needs also to know about the states introduced in CloudStack, especially to those transitional CloudStack VM states like Starting, Migrating, etc.

Upon hypervisor host-connect event, hypervisor resource-agent will first report all VMs on the host to management server, it triggers a "full-sync" process with management server to build an initial sync start point, the host won't be considered as in UP state until this "full-sync" process is completed. After host is connected, hypervisor host resource agent will periodically perform "delta-sync" with CloudStack management server.

 "Full-sync" and "delta-sync" are currently forming the foundation of VMSync process in CloudStack. Although it works nicely most of time for VMs that are solely operated by CloudStack, as soon as the introduction of external VM managers are involved, for example Citrix XenCenter, VMware vCenter, the state sync scenarios can become hard to handle when out-of-band changes posted from external managers, following use cases sometimes can cause problematic issues during normal operations of CloudStack.

1) Takes a long time to bring up all hypervisor hosts in a large setup

During management restart, if things fall out of sync, "full-sync" on host connect-phase can trigger a series of chain actions (actions to bring state in sync) that takes a long time to finish

2) Activities from user, from HA process and VMSync process can collide and the resolution of conflicts is hard to cover all scenarios.

3) Hyprvisor resource-agent to participate into CloudStack VM state management has increased the complexity for people to write a new hypervisor support.

This improvement effort is to address these issues, it will help CloudStack to better interage with third-party virtualization managers like VMware vCenter to perform HA, DRS, FT etc better and reliably through CloudStack.

Design

 High-level principals

At very high-level, we try to attack the problem in following areas

1) Hypervisor resource-agent to report raw VM power state only

This is to de-couple resource agent from CloudStack VM lifecycle state management, letting hypervisor resource-agent only carry on hypervisor-specific actions and report hypervisor raw VM state can greatly simplify the coding of hypervisor resource-agent

In theory, most of currently defined CloudStack transitional VM states are actually representing the states of corresponding transition jobs. For example, CloudStack Starting VM state merely means that there is a pending job in the system which is working on to bring VM from Stopped state to Running state. From end user's perspective, stationary states like Stopped, Running are more meaningful states about a VM. 

There is an important fact that stationary VM states (Stopped, Running) are universal across hypervisors and CloudStack, technically, using stationary VM states and the job status that is currently operating on the subject VM can clearly give user a detail view of the VM. To help us move towards this direction, VM power state is introduced, it currently directly reflects to the hypervisor VM state. With current refactoring work, we still keep original VM state to avoid massive code change and API update, in the future, VM state and VM power state will ultimately be converged into one.

With VM power state, hypervisor resource agent no longer needs to know anything about a transition job status that is specific to CloudStack, all it needs to care is how to carry on a hypervisor-specific action or report VM power state periodically, there is no need to setup a sync start point, therefore we can eliminate "full-sync" process at all.

Following is a code snaplet that shows the old way of how a hypervisor resource agent needs to do a sync report

protected HashMap<String, State> sync() {
        HashMap<String, State> changes = new HashMap<String, State>();
        HashMap<String, State> oldStates = null;

        try {
            synchronized (_vms) {
                HashMap<String, State> newStates = getVmStates();
                oldStates = new HashMap<String, State>(_vms.size());
                oldStates.putAll(_vms);

                for (final Map.Entry<String, State> entry : newStates.entrySet()) {
                    final String vm = entry.getKey();

                    State newState = entry.getValue();
                    final State oldState = oldStates.remove(vm);

                    if (s_logger.isTraceEnabled()) {
                        s_logger.trace("VM " + vm + ": vSphere has state " + newState + " and we have state " + (oldState != null ? oldState.toString() : "null"));
                    }

                    if (vm.startsWith("migrating")) {
                        s_logger.debug("Migrating detected.  Skipping");
                        continue;
                    }

                    if (oldState == null) {
                        _vms.put(vm, newState);
                        s_logger.debug("Detecting a new state but couldn't find a old state so adding it to the changes: " + vm);
                        changes.put(vm, newState);
                    } else if (oldState == State.Starting) {
                        if (newState == State.Running) {
                            _vms.put(vm, newState);
                        } else if (newState == State.Stopped) {
                            s_logger.debug("Ignoring vm " + vm + " because of a lag in starting the vm.");
                        }
                    } else if (oldState == State.Migrating) {
                        if (newState == State.Running) {
                            s_logger.debug("Detected that an migrating VM is now running: " + vm);
                            _vms.put(vm, newState);
                        }
                    } else if (oldState == State.Stopping) {
                        if (newState == State.Stopped) {
                            _vms.put(vm, newState);
                        } else if (newState == State.Running) {
                            s_logger.debug("Ignoring vm " + vm + " because of a lag in stopping the vm. ");
                        }
                    } else if (oldState != newState) {
                        _vms.put(vm, newState);
                        if (newState == State.Stopped) {
                        }
                        changes.put(vm, newState);
                    }
                }

                for (final Map.Entry<String, State> entry : oldStates.entrySet()) {
                    final String vm = entry.getKey();
                    final State oldState = entry.getValue();

                    if (isVmInCluster(vm)) {
                        if (s_logger.isDebugEnabled()) {
                            s_logger.debug("VM " + vm + " is now missing from host report but we detected that it might be migrated to other host by vCenter");
                        }

                        if(oldState != State.Starting && oldState != State.Migrating) {
                            s_logger.debug("VM " + vm + " is now missing from host report and VM is not at starting/migrating state, remove it from host VM-sync map, oldState: " + oldState);
                            _vms.remove(vm);
                        } else {
                            s_logger.debug("VM " + vm + " is missing from host report, but we will ignore VM " + vm + " in transition state " + oldState);
                        }
                        continue;
                    }

                    if (s_logger.isDebugEnabled()) {
                        s_logger.debug("VM " + vm + " is now missing from host report");
                    }

                    if (oldState == State.Stopping) {
                        s_logger.debug("Ignoring VM " + vm + " in transition state stopping.");
                        _vms.remove(vm);
                    } else if (oldState == State.Starting) {
                        s_logger.debug("Ignoring VM " + vm + " in transition state starting.");
                    } else if (oldState == State.Stopped) {
                        _vms.remove(vm);
                    } else if (oldState == State.Migrating) {
                        s_logger.debug("Ignoring VM " + vm + " in migrating state.");
                    } else {
                        State state = State.Stopped;
                        changes.put(entry.getKey(), state);
                    }
                }
            }
        } catch (Throwable e) {
	    // ...
        }
        return changes;
    }

You can see that if CloudStack ever needs to define a new transitional VM state, it will be very hard for hypervisor resource-agent developer to follow, most of sync originated problems happen when developer fails to correctly manage the state cache maintained at resource side.

Since now resource-agent is only required to report raw VM power state, above code logic can become as simple as following

    protected HashMap<String, PowerState> sync() {
    	return getVmStates();
    }

The schema change to add VM power state support is shown below.

ALTER TABLE `cloud`.`vm_instance` ADD COLUMN `power_state` VARCHAR(64) DEFAULT 'PowerUnknown';
ALTER TABLE `cloud`.`vm_instance` ADD COLUMN `power_state_update_time` DATETIME;
ALTER TABLE `cloud`.`vm_instance` ADD COLUMN `power_state_update_count` INT DEFAULT 0;
ALTER TABLE `cloud`.`vm_instance` ADD COLUMN `power_host` bigint unsigned;
ALTER TABLE `cloud`.`vm_instance` ADD CONSTRAINT `fk_vm_instance__power_host` FOREIGN KEY (`power_host`) REFERENCES `cloud`.`host`(`id`);

There is one thing that is worth to note, since majority of time VM will stay in a particular stationary state for a long time, to reduce the number of DB writes, we will only update consecutive same-state update for a limit number of times. power_state_update_count is designed for this purpose.

2) Serialize VM operations

Currently, state transition handling always happens at in-place context, for example, when management server receives hypervisor VM state report, the handling of the report is processed within the context, even if there may be another thread that is handling user request on the same VM. Although we try to coordinate by checking the state of the VM, by simplify failing it with concurrent-access exception.  

About in-place handling style, following code snaplet shows an example.

    
    protected Command compareState(long hostId, VMInstanceVO vm, final AgentVmInfo info, final boolean fullSync, boolean trackExternalChange) {
        State agentState = info.state;
        final State serverState = vm.getState();
        final String serverName = vm.getInstanceName();

        Command command = null;
        s_logger.debug("VM " + serverName + ": cs state = " + serverState + " and realState = " + agentState);
        if (s_logger.isDebugEnabled()) {
            s_logger.debug("VM " + serverName + ": cs state = " + serverState + " and realState = " + agentState);
        }

        if (agentState == State.Error) {
            agentState = State.Stopped;

            short alertType = AlertManager.ALERT_TYPE_USERVM;
            if (VirtualMachine.Type.DomainRouter.equals(vm.getType())) {
                alertType = AlertManager.ALERT_TYPE_DOMAIN_ROUTER;
            } else if (VirtualMachine.Type.ConsoleProxy.equals(vm.getType())) {
                alertType = AlertManager.ALERT_TYPE_CONSOLE_PROXY;
            } else if (VirtualMachine.Type.SecondaryStorageVm.equals(vm.getType())) {
                alertType = AlertManager.ALERT_TYPE_SSVM;
            }

            HostPodVO podVO = _podDao.findById(vm.getPodIdToDeployIn());
            DataCenterVO dcVO = _dcDao.findById(vm.getDataCenterId());
            HostVO hostVO = _hostDao.findById(vm.getHostId());

            String hostDesc = "name: " + hostVO.getName() + " (id:" + hostVO.getId() + "), availability zone: " + dcVO.getName() + ", pod: " + podVO.getName();
            _alertMgr.sendAlert(alertType, vm.getDataCenterId(), vm.getPodIdToDeployIn(), "VM (name: " + vm.getInstanceName() + ", id: " + vm.getId() + ") stopped on host " + hostDesc
                    + " due to storage failure", "Virtual Machine " + vm.getInstanceName() + " (id: " + vm.getId() + ") running on host [" + vm.getHostId() + "] stopped due to storage failure.");
        }

        if (trackExternalChange) {
            if (serverState == State.Starting) {
                if (vm.getHostId() != null && vm.getHostId() != hostId) {
                    s_logger.info("CloudStack is starting VM on host " + vm.getHostId() + ", but status report comes from a different host " + hostId + ", skip status sync for vm: "
                            + vm.getInstanceName());
                    return null;
                }
            }
            if (vm.getHostId() == null || hostId != vm.getHostId()) {
                try {
                    ItWorkVO workItem = _workDao.findByOutstandingWork(vm.getId(), State.Migrating);
                    if(workItem == null){
                        stateTransitTo(vm, VirtualMachine.Event.AgentReportMigrated, hostId);
                    }
                } catch (NoTransitionException e) {
                }
            }
        }

        // during VM migration time, don't sync state will agent status update
        if (serverState == State.Migrating) {
            s_logger.debug("Skipping vm in migrating state: " + vm);
            return null;
        }

        if (trackExternalChange) {
            if (serverState == State.Starting) {
                if (vm.getHostId() != null && vm.getHostId() != hostId) {
                    s_logger.info("CloudStack is starting VM on host " + vm.getHostId() + ", but status report comes from a different host " + hostId + ", skip status sync for vm: "
                            + vm.getInstanceName());
                    return null;
                }
            }

            if (serverState == State.Running) {
                try {
                    //
                    // we had a bug that sometimes VM may be at Running State
                    // but host_id is null, we will cover it here.
                    // means that when CloudStack DB lost of host information,
                    // we will heal it with the info reported from host
                    //
                    if (vm.getHostId() == null || hostId != vm.getHostId()) {
                        if (s_logger.isDebugEnabled()) {
                            s_logger.debug("detected host change when VM " + vm + " is at running state, VM could be live-migrated externally from host " + vm.getHostId() + " to host " + hostId);
                        }

                        stateTransitTo(vm, VirtualMachine.Event.AgentReportMigrated, hostId);
                    }
                } catch (NoTransitionException e) {
                    s_logger.warn(e.getMessage());
                }
            }
        }

        if (agentState == serverState) {
            if (s_logger.isDebugEnabled()) {
                s_logger.debug("Both states are " + agentState + " for " + vm);
            }
            assert (agentState == State.Stopped || agentState == State.Running) : "If the states we send up is changed, this must be changed.";
            if (agentState == State.Running) {
                try {
                    stateTransitTo(vm, VirtualMachine.Event.AgentReportRunning, hostId);
                } catch (NoTransitionException e) {
                    s_logger.warn(e.getMessage());
                }
                // FIXME: What if someone comes in and sets it to stopping? Then
                // what?
                return null;
            }

            s_logger.debug("State matches but the agent said stopped so let's send a cleanup command anyways.");
            return cleanup(vm);
        }

        if (agentState == State.Shutdowned) {
            if (serverState == State.Running || serverState == State.Starting || serverState == State.Stopping) {
                try {
                    advanceStop(vm, true, _accountMgr.getSystemUser(), _accountMgr.getSystemAccount());
                } catch (AgentUnavailableException e) {
                    assert (false) : "How do we hit this with forced on?";
                    return null;
                } catch (OperationTimedoutException e) {
                    assert (false) : "How do we hit this with forced on?";
                    return null;
                } catch (ConcurrentOperationException e) {
                    assert (false) : "How do we hit this with forced on?";
                    return null;
                }
            } else {
                s_logger.debug("Sending cleanup to a shutdowned vm: " + vm.getInstanceName());
                command = cleanup(vm);
            }
        } else if (agentState == State.Stopped) {
            // This state means the VM on the agent was detected previously
            // and now is gone. This is slightly different than if the VM
            // was never completed but we still send down a Stop Command
            // to ensure there's cleanup.
            if (serverState == State.Running) {
                // Our records showed that it should be running so let's restart
                // it.
                _haMgr.scheduleRestart(vm, false);
            } else if (serverState == State.Stopping) {
                _haMgr.scheduleStop(vm, hostId, WorkType.ForceStop);
                s_logger.debug("Scheduling a check stop for VM in stopping mode: " + vm);
            } else if (serverState == State.Starting) {
                s_logger.debug("Ignoring VM in starting mode: " + vm.getInstanceName());
                _haMgr.scheduleRestart(vm, false);
            }
            command = cleanup(vm);
        } else if (agentState == State.Running) {
            if (serverState == State.Starting) {
                if (fullSync) {
                    try {
                        ensureVmRunningContext(hostId, vm, Event.AgentReportRunning);
                    } catch (OperationTimedoutException e) {
                        s_logger.error("Exception during update for running vm: " + vm, e);
                        return null;
                    } catch (ResourceUnavailableException e) {
                        s_logger.error("Exception during update for running vm: " + vm, e);
                        return null;
                    }catch (InsufficientAddressCapacityException e) {
                        s_logger.error("Exception during update for running vm: " + vm, e);
                        return null;
                    }catch (NoTransitionException e) {
                        s_logger.warn(e.getMessage());
                    }
                }
            } else if (serverState == State.Stopping) {
                s_logger.debug("Scheduling a stop command for " + vm);
                _haMgr.scheduleStop(vm, hostId, WorkType.Stop);
            } else {
                s_logger.debug("server VM state " + serverState + " does not meet expectation of a running VM report from agent");

                // just be careful not to stop VM for things we don't handle
                // command = cleanup(vm);
            }
        }
        return command;
    }

In the new design, we will try to serialize activities to the same VM through job facility, since there always be one active operation is in executing, the state transition logic can be simplified. Take the VM migrating case, as it involves with two hosts, in previous model, with VM state report from different hosts, we have to handle it carefully as the host report may come at un-predicted order. 

3) Message bus to coordinate with activities

We will try to use a message-bus to co-ordinate different activities within the management server. This facility is different with the existing feature of "Event Bus", the later one is mainly to integrate external systems through persist-able message-queue servers.

Code changes

1) Message Bus facility

public interface MessageBus {

void setMessageSerializer(MessageSerializer messageSerializer);

MessageSerializer getMessageSerializer();


void subscribe(String topic, MessageSubscriber subscriber);

void unsubscribe(String topic, MessageSubscriber subscriber);

void clearAll();

void prune();


void publish(String senderAddress, String topic, PublishScope scope, Object args);

}

 
MessageBus defines the interface of the message bus facility, it implements a simple publish/subscribe pattern, publishers and subscribers can linked by sharing a common topic,  topic can be in hierarchy mode, a subscriber at higher hierarchy mode can receive messages from all topics that are below. 

MessageBusBase

A simple message bus implementation

MessageHandler

Java annotation for subscriber to specify a message handler

MessageDispatcher

For message subscriber to use to dispatch received messages to annotated message handlers

MessageDetector

To detect interested messages on message bus

2) Job facility

AsyncJobManagerImpl

Refactor it to decouple the tight link with API jobs, make it generic not only executing async API request jobs but also executing internal VM operating jobs

ApiAsyncJobDispatcher

Dispatch async API request jobs

VmWorkJobDispatcher

dispatch internal async VM operation jobs

VmWorkJobVO

VmWorkJobDao

VmWorkJobDaoImpl

Persist classes for internal VM operation jobs

AsyncJobJournalVO

AsyncJobJournalDao

AsyncJobJournalDaoImpl 

Implements job journal facility, all jobs can now have a persist job journal facility

3) VM Power state sync

VirtualMachinePowerStateSync

VirtualMachinePowerStateSyncImpl

4) Other refactored classes

VirtualMachineManagerImpl

HighAvailabilityManagerImpl

VirtualMachineGuru

ReservationContext

Hypervisor resource classes

etc.

Schema changes

ALTER TABLE `cloud`.`async_job` DROP COLUMN `session_key`;

ALTER TABLE `cloud`.`async_job` DROP COLUMN `job_cmd_originator`;

ALTER TABLE `cloud`.`async_job` DROP COLUMN `callback_type`;

ALTER TABLE `cloud`.`async_job` DROP COLUMN `callback_address`;



ALTER TABLE `cloud`.`async_job` ADD COLUMN `parent_id` bigint;

ALTER TABLE `cloud`.`async_job` ADD COLUMN `job_type` VARCHAR(32);

ALTER TABLE `cloud`.`async_job` ADD COLUMN `job_dispatcher` VARCHAR(64);

ALTER TABLE `cloud`.`async_job` ADD COLUMN `job_executing_msid` bigint;



ALTER TABLE `cloud`.`vm_instance` ADD COLUMN `power_state` VARCHAR(74) DEFAULT 'PowerUnknown';

ALTER TABLE `cloud`.`vm_instance` ADD COLUMN `power_state_update_time` DATETIME;

ALTER TABLE `cloud`.`vm_instance` ADD COLUMN `power_state_update_count` INT DEFAULT 0;

ALTER TABLE `cloud`.`vm_instance` ADD COLUMN `power_host` bigint unsigned;

ALTER TABLE `cloud`.`vm_instance` ADD CONSTRAINT `fk_vm_instance__power_host` FOREIGN KEY (`power_host`) REFERENCES `cloud`.`host`(`id`);



CREATE TABLE `cloud`.`vm_work_job` (

  `id` bigint unsigned UNIQUE NOT NULL,

  `step` char(32) NOT NULL COMMENT 'state',

  `vm_type` char(32) NOT NULL COMMENT 'type of vm',

  `vm_instance_id` bigint unsigned NOT NULL COMMENT 'vm instance',

  PRIMARY KEY (`id`),

  CONSTRAINT `fk_vm_work_job__instance_id` FOREIGN KEY (`vm_instance_id`) REFERENCES `vm_instance`(`id`) ON DELETE CASCADE,

  INDEX `i_vm_work_job__vm`(`vm_type`, `vm_instance_id`),

  INDEX `i_vm_work_job__step`(`step`)

) ENGINE=InnoDB DEFAULT CHARSET=utf8;



CREATE TABLE `cloud`.`async_job_journal` (

  `id` bigint unsigned NOT NULL AUTO_INCREMENT COMMENT 'id',

  `job_id` bigint unsigned NOT NULL,

  `journal_type` varchar(32),

  `journal_text` varchar(1024) COMMENT 'journal descriptive informaton',

  `journal_obj` varchar(1024) COMMENT 'journal strutural information, JSON encoded object',

  `created` datetime NOT NULL COMMENT 'date created',

  PRIMARY KEY (`id`),

  CONSTRAINT `fk_async_job_journal__job_id` FOREIGN KEY (`job_id`) REFERENCES `async_job`(`id`) ON DELETE CASCADE

) ENGINE=InnoDB DEFAULT CHARSET=utf8;

 

API/UI

This is low-level change that should keep API compatible, UI change is also not mandatory, we can have UI change to take advantage of better job management in the future(i.e. job journal for more descriptive error messages)