DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.

DUE TO SPAM, SIGN-UP IS DISABLED. Goto Selfserve wiki signup and request an account.
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
...
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.
...
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.
...
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
...
2) Serialize VM operations
Currently, handling of 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.
| Code Block |
|---|
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; } |
...
Follow code snaplet gives a synchronized synchronous handling logic example that is supported in the new model.
| Code Block |
|---|
@Override
public <T extends VMInstanceVO> boolean advanceStop(final T vm, boolean forced, User user, Account account) throws AgentUnavailableException, OperationTimedoutException, ConcurrentOperationException {
VmWorkJobVO workJob = null;
Transaction txn = Transaction.currentTxn();
try {
txn.start();
_vmDao.lockRow(vm.getId(), true);
List<VmWorkJobVO> pendingWorkJobs = _workJobDao.listPendingWorkJobs(
VirtualMachine.Type.Instance, vm.getId(), VmWorkConstants.VM_WORK_STOP);
if(pendingWorkJobs != null && pendingWorkJobs.size() > 0) {
assert(pendingWorkJobs.size() == 1);
workJob = pendingWorkJobs.get(0);
} else {
workJob = new VmWorkJobVO();
workJob.setDispatcher(VmWorkConstants.VM_WORK_JOB_DISPATCHER);
workJob.setCmd(VmWorkConstants.VM_WORK_STOP);
workJob.setAccountId(account.getId());
workJob.setUserId(user.getId());
workJob.setStep(VmWorkJobVO.Step.Prepare);
workJob.setVmType(vm.getType());
workJob.setVmInstanceId(vm.getId());
// save work context info (there are some duplications)
VmWorkStop workInfo = new VmWorkStop();
workInfo.setAccountId(account.getId());
workInfo.setUserId(user.getId());
workInfo.setVmId(vm.getId());
workInfo.setForceStop(forced);
workJob.setCmdInfo(ApiSerializerHelper.toSerializedString(workInfo));
_jobMgr.submitAsyncJob(workJob, VmWorkConstants.VM_WORK_QUEUE, vm.getId());
}
txn.commit();
} catch(Throwable e) {
s_logger.error("Unexpected exception", e);
txn.rollback();
throw new ConcurrentOperationException("Unhandled exception, converted to ConcurrentOperationException");
}
final long jobId = workJob.getId();
AsyncJobExecutionContext.getCurrentExecutionContext().joinJob(jobId);
//
// TODO : this will be replaced with fully-asynchronizedasynchronous way later so that we don't need
// to wait here. The reason we do it synchronizedsynchronous here is that callers of advanceStart is expecting
// synchronizedsynchronous semantics
//
//
_jobMgr.waitAndCheck(
new String[] { TopicConstants.VM_POWER_STATE, TopicConstants.JOB_STATE },
3000L, 600000L, new Predicate() {
@Override
public boolean checkCondition() {
VMInstanceVO instance = _vmDao.findById(vm.getId());
if(instance.getPowerState() == VirtualMachine.PowerState.PowerOff)
return true;
VmWorkJobVO workJob = _workJobDao.findById(jobId);
if(workJob.getStatus() != AsyncJobConstants.STATUS_IN_PROGRESS)
return true;
return false;
}
});
try {
AsyncJobExecutionContext.getCurrentExecutionContext().disjoinJob(jobId);
} catch(Exception e) {
s_logger.error("Unexpected exception", e);
return false;
}
return true;
}
|
...
MessageBus defines the interface of the message bus facility, it implements a simple publish/subscribe pattern, publishers and subscribers can be 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.
...
A job in CloudStack actually represents an orchestration work flow. Due to historic reason, CloudStack has taken considerable efforts trying to make the concept of job implicit to programmers, this is done through the API Command pattern, for API command that is executed asynchronizedlyasynchronously, the request will first be posted to an internal job facility but real execution/processing will be called back into the command object from within the job thread context. The whole job facility has been made implicit intentionally, in most of cases, job facility is used as a context switcher to just provide the execution thread context. This implicit use of job facility actually treats job as secondary class, since explicit job control is discouraged, it leads to the programming model to handle things in-place within the calling context, synchronization is then usually done through locking.
In this refactoring proposal, we will promote jobs into first-class objects, jobs are encouraged to be used in a more explicit way. We will use ordered orderly execution to help reduce the use of locking across the code base and manage orchestration processes explicitly. This will give us better control on managing system load, knows when we need to scale management server cluster and no longer in mercy of java thread pools.
However, due to the legacy bagage we have, moving to this direction will be a long journey, in order to make existing model work without too much change, instead of having one ideal abstract job facility that does general job scheduling, execution and ordering, we will have 3 major job types currently.
...
API job gives a running context for an asynchronized asynchronous API request, it usually starts the an orchestration process.
...
Work job in the new model carries the real orchestrator process, its run will be serialized if related jobs happen to operate on the same underlying target VM object.
3) Pseudo Job
Inside CloudStack, there are a few manager components that use their own threads to manage service activities, when it comes to use the newly introduced work jobs for orchestration, we sometimes need a pseudo job context, pseudo job provides just that context. The difference between Pseudo job and an high level API job is that pseudo job runs in its own thread context, while API job runs in the thread from job thread pool.
...
Like a process in operating system, a job can have multiple execution states, it could be put in blocking, or be in currently running state, etc. Joining another job means to wait for completion of the subject job, be either a successful completion or a failure completion. A blocked job may be rescheduled to run based on triggering of events. Currently, a job that is joining to another thread job can only be scheduled to run on upon wakeup events.
Job wakeup
When a job joins to another job, to wait for the completion status of joined job, there are two ways to achieve that. We've shown it for the first way in _jobMgr.waitAndCheck - blocking the executing thread until the condition is satisfied. The problem of this approach is that it holds an a real executing thread. If a caller already has a persistent thread, it is not a problem, however, for most of API initiated orchestration jobs, they all share a global job thread pool, blocking executing thread is not the most efficient way for system scalability. The new job facility provides a support for a second approach, this approach will put the job into blocking state, release the executing thread, and then reschedule job execution based on wakeup calls(event triggered).
Message bus provides delivery-at-best service to CloudStack components, for locally broadcast messages(within one management server), it is reliable when management server is running, however, for messages that are across-ing management server boundaries, it is not a 100 percent reliable service for building a reliable orchestration process. When a job is joining and waiting for another job to complete, in order for the job check-up process keep - on going, the job will be periodically waken up on a specified interval. Therefore, message bus service can be used for efficient event notification and in case that message bus service fails, this wakeup service will help ensure the reliability of the whole orchestration process.
...
When there is no pending job working on the VM, VM should always stay at stationary states (i.e., PoweredOn or PoweredOff), out of sync situation between what CloudStack DB has record and what a host has reported will be resolved with a new serialized job flow, depends on HA configuration, we can either try to eventually bring VM state to be in sync with CloudStack DB or let CloudStack honor what it is reported.
...
| Code Block |
|---|
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 `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`.`async_job` ADD COLUMN `job_pending_signals` int(10) NOT NULL DEFAULT 0; 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`); 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; CREATE TABLE `cloud`.`async_job_join_map` ( `id` bigint unsigned NOT NULL AUTO_INCREMENT COMMENT 'id', `job_id` bigint unsigned NOT NULL, `join_job_id` bigint unsigned NOT NULL, `join_status` int NOT NULL, `join_result` varchar(1024), `join_msid` bigint, `complete_msid` bigint, `sync_source_id` bigint COMMENT 'upper-level job sync source info before join', `wakeup_handler` varchar(64), `wakeup_dispatcher` varchar(64), `wakeup_interval` bigint NOT NULL DEFAULT 3000 COMMENT 'wakeup interval in seconds', `created` datetime NOT NULL, `last_updated` datetime, `next_wakeup` datetime, `expiration` datetime, PRIMARY KEY (`id`), CONSTRAINT `fk_async_job_join_map__job_id` FOREIGN KEY (`job_id`) REFERENCES `async_job`(`id`) ON DELETE CASCADE, CONSTRAINT `fk_async_job_join_map__join_job_id` FOREIGN KEY (`join_job_id`) REFERENCES `async_job`(`id`), CONSTRAINT `fk_async_job_join_map__join` UNIQUE (`job_id`, `join_job_id`), INDEX `i_async_job_join_map__join_job_id`(`join_job_id`), INDEX `i_async_job_join_map__created`(`created`), INDEX `i_async_job_join_map__last_updated`(`last_updated`), INDEX `i_async_job_join_map__next_wakeup`(`next_wakeup`), INDEX `i_async_job_join_map__expiration`(`expiration`) ) ENGINE=InnoDB DEFAULT CHARSET=utf8; |
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)