Skip to content

Commit 3574d8d

Browse files
DaanHooglandDaan HooglandWei Zhou
authored
parallel nic adding (apache#5541)
* trace nics additions * work queue patch for network to add * add secondary key to job * logging improvements and naming of field(s) * several naming corrections * extra check if net already exists for vm * placeholder job with secondary object * constraint on entering the same job multiple times * error handling/warning message * review comments applied Co-authored-by: Daan Hoogland <[email protected]> Co-authored-by: Wei Zhou <[email protected]>
1 parent 93c0b60 commit 3574d8d

7 files changed

Lines changed: 127 additions & 29 deletions

File tree

‎engine/orchestration/src/main/java/com/cloud/vm/VirtualMachineManagerImpl.java‎

Lines changed: 68 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -42,6 +42,7 @@
4242

4343
import javax.inject.Inject;
4444
import javax.naming.ConfigurationException;
45+
import javax.persistence.EntityExistsException;
4546

4647
import org.apache.cloudstack.affinity.dao.AffinityGroupVMMapDao;
4748
import org.apache.cloudstack.annotation.AnnotationService;
@@ -3981,7 +3982,7 @@ public NicProfile addVmToNetwork(final VirtualMachine vm, final Network network,
39813982
if (jobContext.isJobDispatchedBy(VmWorkConstants.VM_WORK_JOB_DISPATCHER)) {
39823983
// avoid re-entrance
39833984
VmWorkJobVO placeHolder = null;
3984-
placeHolder = createPlaceHolderWork(vm.getId());
3985+
placeHolder = createPlaceHolderWork(vm.getId(), network.getUuid());
39853986
try {
39863987
return orchestrateAddVmToNetwork(vm, network, requested);
39873988
} finally {
@@ -4021,10 +4022,23 @@ public NicProfile addVmToNetwork(final VirtualMachine vm, final Network network,
40214022
}
40224023
}
40234024

4025+
/**
4026+
* duplicated in {@see UserVmManagerImpl} for a {@see UserVmVO}
4027+
*/
4028+
private void checkIfNetworkExistsForVM(VirtualMachine virtualMachine, Network network) {
4029+
List<NicVO> allNics = _nicsDao.listByVmId(virtualMachine.getId());
4030+
for (NicVO nic : allNics) {
4031+
if (nic.getNetworkId() == network.getId()) {
4032+
throw new CloudRuntimeException("A NIC already exists for VM:" + virtualMachine.getInstanceName() + " in network: " + network.getUuid());
4033+
}
4034+
}
4035+
}
4036+
40244037
private NicProfile orchestrateAddVmToNetwork(final VirtualMachine vm, final Network network, final NicProfile requested) throws ConcurrentOperationException, ResourceUnavailableException,
40254038
InsufficientCapacityException {
40264039
final CallContext cctx = CallContext.current();
40274040

4041+
checkIfNetworkExistsForVM(vm, network);
40284042
s_logger.debug("Adding vm " + vm + " to network " + network + "; requested nic profile " + requested);
40294043
final VMInstanceVO vmVO = _vmDao.findById(vm.getId());
40304044
final ReservationContext context = new ReservationContextImpl(null, null, cctx.getCallingUser(), cctx.getCallingAccount());
@@ -5375,7 +5389,7 @@ public Outcome<VirtualMachine> migrateVmThroughJobQueue(final String vmUuid, fin
53755389
Map<Volume, StoragePool> volumeStorageMap = dest.getStorageForDisks();
53765390
if (volumeStorageMap != null) {
53775391
for (Volume vol : volumeStorageMap.keySet()) {
5378-
checkConcurrentJobsPerDatastoreThreshhold(volumeStorageMap.get(vol));
5392+
checkConcurrentJobsPerDatastoreThreshold(volumeStorageMap.get(vol));
53795393
}
53805394
}
53815395

@@ -5540,7 +5554,7 @@ public Outcome<VirtualMachine> migrateVmForScaleThroughJobQueue(
55405554
return new VmJobVirtualMachineOutcome(workJob, vm.getId());
55415555
}
55425556

5543-
private void checkConcurrentJobsPerDatastoreThreshhold(final StoragePool destPool) {
5557+
private void checkConcurrentJobsPerDatastoreThreshold(final StoragePool destPool) {
55445558
final Long threshold = VolumeApiService.ConcurrentMigrationsThresholdPerDatastore.value();
55455559
if (threshold != null && threshold > 0) {
55465560
long count = _jobMgr.countPendingJobs("\"storageid\":\"" + destPool.getUuid() + "\"", MigrateVMCmd.class.getName(), MigrateVolumeCmd.class.getName(), MigrateVolumeCmdByAdmin.class.getName());
@@ -5561,7 +5575,7 @@ public Outcome<VirtualMachine> migrateVmStorageThroughJobQueue(
55615575
Set<Long> uniquePoolIds = new HashSet<>(poolIds);
55625576
for (Long poolId : uniquePoolIds) {
55635577
StoragePoolVO pool = _storagePoolDao.findById(poolId);
5564-
checkConcurrentJobsPerDatastoreThreshhold(pool);
5578+
checkConcurrentJobsPerDatastoreThreshold(pool);
55655579
}
55665580

55675581
final VMInstanceVO vm = _vmDao.findByUuid(vmUuid);
@@ -5608,35 +5622,61 @@ public Outcome<VirtualMachine> addVmToNetworkThroughJobQueue(
56085622

56095623
final List<VmWorkJobVO> pendingWorkJobs = _workJobDao.listPendingWorkJobs(
56105624
VirtualMachine.Type.Instance, vm.getId(),
5611-
VmWorkAddVmToNetwork.class.getName());
5625+
VmWorkAddVmToNetwork.class.getName(), network.getUuid());
56125626

56135627
VmWorkJobVO workJob = null;
56145628
if (pendingWorkJobs != null && pendingWorkJobs.size() > 0) {
5615-
assert pendingWorkJobs.size() == 1;
5629+
if (pendingWorkJobs.size() > 1) {
5630+
s_logger.warn(String.format("The number of jobs to add network %s to vm %s are %d", network.getUuid(), vm.getInstanceName(), pendingWorkJobs.size()));
5631+
}
56165632
workJob = pendingWorkJobs.get(0);
56175633
} else {
5634+
if (s_logger.isTraceEnabled()) {
5635+
s_logger.trace(String.format("no jobs to add network %s for vm %s yet", network, vm));
5636+
}
56185637

5619-
workJob = new VmWorkJobVO(context.getContextId());
5638+
workJob = createVmWorkJobToAddNetwork(vm, network, requested, context, user, account);
5639+
}
5640+
AsyncJobExecutionContext.getCurrentExecutionContext().joinJob(workJob.getId());
56205641

5621-
workJob.setDispatcher(VmWorkConstants.VM_WORK_JOB_DISPATCHER);
5622-
workJob.setCmd(VmWorkAddVmToNetwork.class.getName());
5642+
return new VmJobVirtualMachineOutcome(workJob, vm.getId());
5643+
}
56235644

5624-
workJob.setAccountId(account.getId());
5625-
workJob.setUserId(user.getId());
5626-
workJob.setVmType(VirtualMachine.Type.Instance);
5627-
workJob.setVmInstanceId(vm.getId());
5628-
workJob.setRelated(AsyncJobExecutionContext.getOriginJobId());
5645+
private VmWorkJobVO createVmWorkJobToAddNetwork(
5646+
VirtualMachine vm,
5647+
Network network,
5648+
NicProfile requested,
5649+
CallContext context,
5650+
User user,
5651+
Account account) {
5652+
VmWorkJobVO workJob;
5653+
workJob = new VmWorkJobVO(context.getContextId());
56295654

5630-
// save work context info (there are some duplications)
5631-
final VmWorkAddVmToNetwork workInfo = new VmWorkAddVmToNetwork(user.getId(), account.getId(), vm.getId(),
5632-
VirtualMachineManagerImpl.VM_WORK_JOB_HANDLER, network.getId(), requested);
5633-
workJob.setCmdInfo(VmWorkSerializer.serialize(workInfo));
5655+
workJob.setDispatcher(VmWorkConstants.VM_WORK_JOB_DISPATCHER);
5656+
workJob.setCmd(VmWorkAddVmToNetwork.class.getName());
56345657

5658+
workJob.setAccountId(account.getId());
5659+
workJob.setUserId(user.getId());
5660+
workJob.setVmType(VirtualMachine.Type.Instance);
5661+
workJob.setVmInstanceId(vm.getId());
5662+
workJob.setRelated(AsyncJobExecutionContext.getOriginJobId());
5663+
workJob.setSecondaryObjectIdentifier(network.getUuid());
5664+
5665+
// save work context info (there are some duplications)
5666+
final VmWorkAddVmToNetwork workInfo = new VmWorkAddVmToNetwork(user.getId(), account.getId(), vm.getId(),
5667+
VirtualMachineManagerImpl.VM_WORK_JOB_HANDLER, network.getId(), requested);
5668+
workJob.setCmdInfo(VmWorkSerializer.serialize(workInfo));
5669+
5670+
try {
56355671
_jobMgr.submitAsyncJob(workJob, VmWorkConstants.VM_WORK_QUEUE, vm.getId());
5672+
} catch (CloudRuntimeException e) {
5673+
if (e.getCause() instanceof EntityExistsException) {
5674+
String msg = String.format("A job to add a nic for network %s to vm %s already exists", network.getUuid(), vm.getUuid());
5675+
s_logger.warn(msg, e);
5676+
}
5677+
throw e;
56365678
}
5637-
AsyncJobExecutionContext.getCurrentExecutionContext().joinJob(workJob.getId());
5638-
5639-
return new VmJobVirtualMachineOutcome(workJob, vm.getId());
5679+
return workJob;
56405680
}
56415681

56425682
public Outcome<VirtualMachine> removeNicFromVmThroughJobQueue(
@@ -5945,6 +5985,10 @@ public Pair<JobInfo.Status, String> handleVmWorkJob(final VmWork work) throws Ex
59455985
}
59465986

59475987
private VmWorkJobVO createPlaceHolderWork(final long instanceId) {
5988+
return createPlaceHolderWork(instanceId, null);
5989+
}
5990+
5991+
private VmWorkJobVO createPlaceHolderWork(final long instanceId, String secondaryObjectIdentifier) {
59485992
final VmWorkJobVO workJob = new VmWorkJobVO("");
59495993

59505994
workJob.setDispatcher(VmWorkConstants.VM_WORK_JOB_PLACEHOLDER);
@@ -5956,6 +6000,9 @@ private VmWorkJobVO createPlaceHolderWork(final long instanceId) {
59566000
workJob.setStep(VmWorkJobVO.Step.Starting);
59576001
workJob.setVmType(VirtualMachine.Type.Instance);
59586002
workJob.setVmInstanceId(instanceId);
6003+
if(StringUtils.isNotBlank(secondaryObjectIdentifier)) {
6004+
workJob.setSecondaryObjectIdentifier(secondaryObjectIdentifier);
6005+
}
59596006
workJob.setInitMsid(ManagementServerNode.getManagementServerId());
59606007

59616008
_workJobDao.persist(workJob);

‎engine/schema/src/main/resources/META-INF/db/schema-41520to41600.sql‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -791,3 +791,6 @@ ALTER TABLE cloud.user_vm_details MODIFY value varchar(5120) NOT NULL;
791791
ALTER TABLE cloud_usage.usage_network DROP PRIMARY KEY, ADD PRIMARY KEY (`account_id`,`zone_id`,`host_id`,`network_id`,`event_time_millis`);
792792
ALTER TABLE `cloud`.`user_statistics` DROP INDEX `account_id`, ADD UNIQUE KEY `account_id` (`account_id`,`data_center_id`,`public_ip_address`,`device_id`,`device_type`, `network_id`);
793793
ALTER TABLE `cloud_usage`.`user_statistics` DROP INDEX `account_id`, ADD UNIQUE KEY `account_id` (`account_id`,`data_center_id`,`public_ip_address`,`device_id`,`device_type`, `network_id`);
794+
795+
ALTER TABLE `cloud`.`vm_work_job` ADD COLUMN `secondary_object` char(100) COMMENT 'any additional item that must be checked during queueing' AFTER `vm_instance_id`;
796+
ALTER TABLE cloud.vm_work_job ADD CONSTRAINT vm_work_job_step_and_objects UNIQUE KEY (step,vm_instance_id,secondary_object);

‎framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/dao/VmWorkJobDao.java‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,8 @@ public interface VmWorkJobDao extends GenericDao<VmWorkJobVO, Long> {
3232

3333
List<VmWorkJobVO> listPendingWorkJobs(VirtualMachine.Type type, long instanceId, String jobCmd);
3434

35+
List<VmWorkJobVO> listPendingWorkJobs(VirtualMachine.Type type, long instanceId, String jobCmd, String secondaryObjectIdentifier);
36+
3537
void updateStep(long workJobId, Step step);
3638

3739
void expungeCompletedWorkJobs(Date cutDate);

‎framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/dao/VmWorkJobDaoImpl.java‎

Lines changed: 15 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -67,6 +67,7 @@ public void init() {
6767
PendingWorkJobByCommandSearch.and("jobStatus", PendingWorkJobByCommandSearch.entity().getStatus(), Op.EQ);
6868
PendingWorkJobByCommandSearch.and("vmType", PendingWorkJobByCommandSearch.entity().getVmType(), Op.EQ);
6969
PendingWorkJobByCommandSearch.and("vmInstanceId", PendingWorkJobByCommandSearch.entity().getVmInstanceId(), Op.EQ);
70+
PendingWorkJobByCommandSearch.and("secondaryObjectIdentifier", PendingWorkJobByCommandSearch.entity().getSecondaryObjectIdentifier(), Op.EQ);
7071
PendingWorkJobByCommandSearch.and("step", PendingWorkJobByCommandSearch.entity().getStep(), Op.NEQ);
7172
PendingWorkJobByCommandSearch.and("cmd", PendingWorkJobByCommandSearch.entity().getCmd(), Op.EQ);
7273
PendingWorkJobByCommandSearch.done();
@@ -119,6 +120,20 @@ public List<VmWorkJobVO> listPendingWorkJobs(VirtualMachine.Type type, long inst
119120
return this.listBy(sc, filter);
120121
}
121122

123+
@Override
124+
public List<VmWorkJobVO> listPendingWorkJobs(VirtualMachine.Type type, long instanceId, String jobCmd, String secondaryObjectIdentifier) {
125+
126+
SearchCriteria<VmWorkJobVO> sc = PendingWorkJobByCommandSearch.create();
127+
sc.setParameters("jobStatus", JobInfo.Status.IN_PROGRESS);
128+
sc.setParameters("vmType", type);
129+
sc.setParameters("vmInstanceId", instanceId);
130+
sc.setParameters("secondaryObjectIdentifier", secondaryObjectIdentifier);
131+
sc.setParameters("cmd", jobCmd);
132+
133+
Filter filter = new Filter(VmWorkJobVO.class, "created", true, null, null);
134+
return this.listBy(sc, filter);
135+
}
136+
122137
@Override
123138
public void updateStep(long workJobId, Step step) {
124139
VmWorkJobVO jobVo = findById(workJobId);

‎framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/impl/AsyncJobVO.java‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -384,7 +384,7 @@ public void setRemoved(final Date removed) {
384384
@Override
385385
public String toString() {
386386
StringBuffer sb = new StringBuffer();
387-
sb.append("AsyncJobVO {id:").append(getId());
387+
sb.append("AsyncJobVO : {id:").append(getId());
388388
sb.append(", userId: ").append(getUserId());
389389
sb.append(", accountId: ").append(getAccountId());
390390
sb.append(", instanceType: ").append(getInstanceType());

‎framework/jobs/src/main/java/org/apache/cloudstack/framework/jobs/impl/VmWorkJobVO.java‎

Lines changed: 24 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -58,6 +58,9 @@ boolean updateState() {
5858
@Column(name = "vm_instance_id")
5959
long vmInstanceId;
6060

61+
@Column(name = "secondary_object")
62+
String secondaryObjectIdentifier;
63+
6164
protected VmWorkJobVO() {
6265
}
6366

@@ -89,4 +92,25 @@ public long getVmInstanceId() {
8992
public void setVmInstanceId(long vmInstanceId) {
9093
this.vmInstanceId = vmInstanceId;
9194
}
95+
96+
public String getSecondaryObjectIdentifier() {
97+
return secondaryObjectIdentifier;
98+
}
99+
100+
public void setSecondaryObjectIdentifier(String secondaryObjectIdentifier) {
101+
this.secondaryObjectIdentifier = secondaryObjectIdentifier;
102+
}
103+
@Override
104+
public String toString() {
105+
StringBuffer sb = new StringBuffer();
106+
sb.append("VmWorkJobVO : {").
107+
append(", step: ").append(getStep()).
108+
append(", vmType: ").append(getVmType()).
109+
append(", vmInstanceId: ").append(getVmInstanceId()).
110+
append(", secondaryObjectIdentifier: ").append(getSecondaryObjectIdentifier()).
111+
append(super.toString()).
112+
append("}");
113+
return sb.toString();
114+
}
115+
92116
}

‎server/src/main/java/com/cloud/vm/UserVmManagerImpl.java‎

Lines changed: 14 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -1384,12 +1384,7 @@ public UserVm addNicToVirtualMachine(AddNicToVMCmd cmd) throws InvalidParameterV
13841384
Account vmOwner = _accountMgr.getAccount(vmInstance.getAccountId());
13851385
_networkModel.checkNetworkPermissions(vmOwner, network);
13861386

1387-
List<NicVO> allNics = _nicDao.listByVmId(vmInstance.getId());
1388-
for (NicVO nic : allNics) {
1389-
if (nic.getNetworkId() == network.getId()) {
1390-
throw new CloudRuntimeException("A NIC already exists for VM:" + vmInstance.getInstanceName() + " in network: " + network.getUuid());
1391-
}
1392-
}
1387+
checkIfNetExistsForVM(vmInstance, network);
13931388

13941389
macAddress = validateOrReplaceMacAddress(macAddress, network.getId());
13951390

@@ -1456,10 +1451,22 @@ public UserVm addNicToVirtualMachine(AddNicToVMCmd cmd) throws InvalidParameterV
14561451
}
14571452
}
14581453
CallContext.current().putContextParameter(Nic.class, guestNic.getUuid());
1459-
s_logger.debug("Successful addition of " + network + " from " + vmInstance);
1454+
s_logger.debug(String.format("Successful addition of %s from %s through %s", network, vmInstance, guestNic));
14601455
return _vmDao.findById(vmInstance.getId());
14611456
}
14621457

1458+
/**
1459+
* duplicated in {@see VirtualMachineManagerImpl} for a {@see VMInstanceVO}
1460+
*/
1461+
private void checkIfNetExistsForVM(VirtualMachine virtualMachine, Network network) {
1462+
List<NicVO> allNics = _nicDao.listByVmId(virtualMachine.getId());
1463+
for (NicVO nic : allNics) {
1464+
if (nic.getNetworkId() == network.getId()) {
1465+
throw new CloudRuntimeException("A NIC already exists for VM:" + virtualMachine.getInstanceName() + " in network: " + network.getUuid());
1466+
}
1467+
}
1468+
}
1469+
14631470
/**
14641471
* If the given MAC address is invalid it replaces the given MAC with the next available MAC address
14651472
*/

0 commit comments

Comments
 (0)