diff --git a/plugin/flatNetworkProvider/src/main/java/org/zstack/network/service/flat/FlatEipBackend.java b/plugin/flatNetworkProvider/src/main/java/org/zstack/network/service/flat/FlatEipBackend.java index 0500fb4e2ea..633c7f0489a 100755 --- a/plugin/flatNetworkProvider/src/main/java/org/zstack/network/service/flat/FlatEipBackend.java +++ b/plugin/flatNetworkProvider/src/main/java/org/zstack/network/service/flat/FlatEipBackend.java @@ -13,6 +13,7 @@ import org.zstack.core.errorcode.ErrorFacade; import org.zstack.core.timeout.ApiTimeoutManager; import org.zstack.header.core.Completion; +import org.zstack.header.core.FutureCompletion; import org.zstack.header.core.NopeCompletion; import org.zstack.header.core.workflow.Flow; import org.zstack.header.core.workflow.FlowRollback; @@ -57,6 +58,7 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import static java.util.Arrays.asList; @@ -120,6 +122,7 @@ public static class DeleteEipCmd extends AgentCmd { public static class BatchApplyEipCmd extends AgentCmd { public List eips; + public boolean prepare; } public static class BatchDeleteEipCmd extends AgentCmd { @@ -133,6 +136,17 @@ public static class BatchDeleteEipCmd extends AgentCmd { @Override public void preMigrateVm(VmInstanceInventory inv, String destHostUuid) { + List eips = getEipsByVmUuid(inv.getUuid()); + if (eips == null || eips.isEmpty()) { + return; + } + + FutureCompletion completion = new FutureCompletion(null); + batchApplyEips(eips, destHostUuid, true, false, completion); + completion.await(TimeUnit.MINUTES.toMillis(30)); + if (!completion.isSuccess()) { + throw new OperationFailureException(completion.getErrorCode()); + } } @Override @@ -147,46 +161,54 @@ public void afterMigrateVm(final VmInstanceInventory inv, String srcHostUuid) { return; } - batchDeleteEips(eips, srcHostUuid, new Completion(null) { + batchApplyEips(eips, inv.getHostUuid(), new Completion(null) { @Override public void success() { - batchApplyEips(eips, inv.getHostUuid(), new Completion(null) { + batchDeleteEips(eips, srcHostUuid, new Completion(null) { @Override public void success() { - logger.warn(String.format("after migration, successfully applied EIPs[uuids:%s] to the vm[uuid:%s, name:%s] on the destination host[uuid:%s] after delete eip on src host[uuid:%s] succeeded", - eips.stream().map(e -> e.vip).collect(Collectors.toList()), inv.getUuid(), inv.getName(), inv.getHostUuid(), srcHostUuid)); } @Override public void fail(ErrorCode errorCode) { - logger.warn(String.format("after migration, failed to apply EIPs[uuids:%s] to the vm[uuid:%s, name:%s] on the destination host[uuid:%s] after delete eip on src host[uuid:%s] succeeded, %s", - eips.stream().map(e -> e.vip).collect(Collectors.toList()), inv.getUuid(), inv.getName(), inv.getHostUuid(), srcHostUuid, errorCode)); + logger.warn(String.format("failed to delete EIPs[vips:%s] for migrated vm[uuid:%s] on source host[uuid:%s], %s", + eips.stream().map(e -> e.vip).collect(Collectors.toList()), inv.getUuid(), srcHostUuid, errorCode)); } }); } @Override public void fail(ErrorCode errorCode) { - batchApplyEips(eips, inv.getHostUuid(), new Completion(null) { - @Override - public void success() { - logger.warn(String.format("after migration, successfully applied EIPs[uuids:%s] to the vm[uuid:%s, name:%s] on the destination host[uuid:%s] after delete eip on src host[uuid:%s] failed", - eips.stream().map(e -> e.vip).collect(Collectors.toList()), inv.getUuid(), inv.getName(), inv.getHostUuid(), srcHostUuid)); - } - - @Override - public void fail(ErrorCode errorCode) { - logger.warn(String.format("after migration, failed to apply EIPs[uuids:%s] to the vm[uuid:%s, name:%s] on the destination host[uuid:%s] after delete eip on src host[uuid:%s] failed, %s", - eips.stream().map(e -> e.vip).collect(Collectors.toList()), inv.getUuid(), inv.getName(), inv.getHostUuid(), srcHostUuid, errorCode)); - } - }); + logger.warn(String.format("failed to enable EIPs[vips:%s] for migrated vm[uuid:%s] on destination host[uuid:%s], keep source EIPs on host[uuid:%s], %s", + eips.stream().map(e -> e.vip).collect(Collectors.toList()), inv.getUuid(), inv.getHostUuid(), srcHostUuid, errorCode)); } }); } @Override public void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason) { + List eips = getEipsByVmUuid(inv.getUuid()); + if (eips == null || eips.isEmpty() || destHostUuid == null) { + return; + } + if (HostErrors.FAILED_TO_MIGRATE_VM_ON_HYPERVISOR.isEqual(reason.getCode())) { + logger.warn(String.format("keep prepared EIPs[vips:%s] on destination host[uuid:%s] because vm[uuid:%s] placement is uncertain after migration failure", + eips.stream().map(e -> e.vip).collect(Collectors.toList()), destHostUuid, inv.getUuid())); + return; + } + + batchDeleteEips(eips, destHostUuid, new Completion(null) { + @Override + public void success() { + } + + @Override + public void fail(ErrorCode errorCode) { + logger.warn(String.format("failed to clean prepared EIPs[vips:%s] for vm[uuid:%s] on destination host[uuid:%s], %s", + eips.stream().map(e -> e.vip).collect(Collectors.toList()), inv.getUuid(), destHostUuid, errorCode)); + } + }); } @Transactional(readOnly = true) @@ -250,13 +272,25 @@ public void fail(ErrorCode errorCode) { } private void vmMigrateToAnotherHost(final FlowTrigger trigger) { - batchDeleteEips(eips, struct.getOriginalHostUuid(), new NopeCompletion()); - applyHostUuidForRollback = struct.getOriginalHostUuid(); batchApplyEips(eips, struct.getCurrentHostUuid(), new Completion(trigger) { @Override public void success() { releaseHostUuidForRollback = struct.getCurrentHostUuid(); - trigger.next(); + batchDeleteEips(eips, struct.getOriginalHostUuid(), new Completion(trigger) { + @Override + public void success() { + applyHostUuidForRollback = struct.getOriginalHostUuid(); + trigger.next(); + } + + @Override + public void fail(ErrorCode errorCode) { + releaseHostUuidForRollback = null; + logger.warn(String.format("failed to clean EIPs[vips:%s] for vm[uuid:%s] on original host[uuid:%s] after applying them on current host[uuid:%s], %s", + eips.stream().map(e -> e.vip).collect(Collectors.toList()), vm.getUuid(), struct.getOriginalHostUuid(), struct.getCurrentHostUuid(), errorCode)); + trigger.next(); + } + }); } @Override @@ -268,13 +302,25 @@ public void fail(ErrorCode errorCode) { private void vmRunningFromUnknownStateHostChanged(final FlowTrigger trigger) { - batchDeleteEips(eips, struct.getOriginalHostUuid(), new NopeCompletion()); - applyHostUuidForRollback = struct.getOriginalHostUuid(); batchApplyEips(eips, struct.getCurrentHostUuid(), new Completion(trigger) { @Override public void success() { releaseHostUuidForRollback = struct.getCurrentHostUuid(); - trigger.next(); + batchDeleteEips(eips, struct.getOriginalHostUuid(), new Completion(trigger) { + @Override + public void success() { + applyHostUuidForRollback = struct.getOriginalHostUuid(); + trigger.next(); + } + + @Override + public void fail(ErrorCode errorCode) { + releaseHostUuidForRollback = null; + logger.warn(String.format("failed to clean EIPs[vips:%s] for vm[uuid:%s] on original host[uuid:%s] after applying them on current host[uuid:%s], %s", + eips.stream().map(e -> e.vip).collect(Collectors.toList()), vm.getUuid(), struct.getOriginalHostUuid(), struct.getCurrentHostUuid(), errorCode)); + trigger.next(); + } + }); } @Override @@ -519,12 +565,18 @@ public void run(MessageReply reply) { } private void batchApplyEips(List eips, String hostUuid, final Completion completion) { - batchApplyEips(eips, hostUuid, false, completion); + batchApplyEips(eips, hostUuid, false, false, completion); } private void batchApplyEips(List eips, String hostUuid, boolean noHostStatusCheck, final Completion completion) { + batchApplyEips(eips, hostUuid, false, noHostStatusCheck, completion); + } + + private void batchApplyEips(List eips, String hostUuid, boolean prepare, boolean noHostStatusCheck, + final Completion completion) { BatchApplyEipCmd cmd = new BatchApplyEipCmd(); cmd.eips = eips; + cmd.prepare = prepare; KVMHostAsyncHttpCallMsg msg = new KVMHostAsyncHttpCallMsg(); msg.setCommand(cmd); @@ -547,6 +599,11 @@ public void run(MessageReply reply) { return; } + if (prepare) { + completion.success(); + return; + } + List vipUuids = CollectionUtils.transformToList(eips, new Function() { @Override public String call(EipTO arg) { diff --git a/test/src/test/groovy/org/zstack/test/integration/networkservice/provider/flat/eip/StartFlatNetworkVmWithEipCase.groovy b/test/src/test/groovy/org/zstack/test/integration/networkservice/provider/flat/eip/StartFlatNetworkVmWithEipCase.groovy index ac5cfa42337..68477e79e49 100644 --- a/test/src/test/groovy/org/zstack/test/integration/networkservice/provider/flat/eip/StartFlatNetworkVmWithEipCase.groovy +++ b/test/src/test/groovy/org/zstack/test/integration/networkservice/provider/flat/eip/StartFlatNetworkVmWithEipCase.groovy @@ -2,7 +2,16 @@ package org.zstack.test.integration.networkservice.provider.flat.eip import org.springframework.http.HttpEntity import org.zstack.core.cloudbus.CloudBus +import org.zstack.core.workflow.FlowChainBuilder +import org.zstack.header.core.FutureCompletion +import org.zstack.header.core.workflow.Flow +import org.zstack.header.core.workflow.FlowDoneHandler +import org.zstack.header.core.workflow.FlowErrorHandler +import org.zstack.header.core.workflow.FlowRollback +import org.zstack.header.core.workflow.FlowTrigger import org.zstack.header.network.service.NetworkServiceType +import org.zstack.header.vm.VmAbnormalLifeCycleStruct +import org.zstack.header.vm.VmInstanceVO import org.zstack.network.service.eip.EipConstant import org.zstack.network.service.flat.FlatEipBackend import org.zstack.network.service.flat.FlatNetworkServiceConstant @@ -18,6 +27,10 @@ import org.zstack.testlib.HttpError import org.zstack.testlib.SubCase import org.zstack.utils.data.SizeUnit +import java.util.concurrent.TimeUnit +import java.util.concurrent.atomic.AtomicBoolean +import java.util.concurrent.atomic.AtomicInteger + import static org.zstack.core.Platform.operr /** @@ -245,15 +258,17 @@ class StartFlatNetworkVmWithEipCase extends SubCase { vmNicUuid = vm.vmNics[0].uuid } - boolean deleteEipOnSrcHostFailed = false + testAbnormalMigrationSourceCleanupFailure() + + List migrationEipOperations = Collections.synchronizedList([]) env.afterSimulator(FlatEipBackend.BATCH_DELETE_EIP_PATH) { - deleteEipOnSrcHostFailed = true + migrationEipOperations.add("delete") throw new HttpError(403, "on purpose") } - boolean applyEipOnDestHostSuccessed = false env.afterSimulator(FlatEipBackend.BATCH_APPLY_EIP_PATH) { rsp, HttpEntity e -> - applyEipOnDestHostSuccessed = true + def cmd = json(e.body, FlatEipBackend.BatchApplyEipCmd.class) + migrationEipOperations.add(cmd.prepare ? "prepare" : "active") return rsp } @@ -263,13 +278,17 @@ class StartFlatNetworkVmWithEipCase extends SubCase { } retryInSecs { - assert deleteEipOnSrcHostFailed == true - assert applyEipOnDestHostSuccessed == true + int prepareIndex = migrationEipOperations.indexOf("prepare") + int activeIndex = migrationEipOperations.indexOf("active") + int deleteIndex = migrationEipOperations.indexOf("delete") + assert prepareIndex >= 0 : "EIP migration did not prepare destination EIP: operations=${migrationEipOperations}" + assert activeIndex > prepareIndex : "Destination EIP was not activated after prepare: operations=${migrationEipOperations}" + assert deleteIndex > activeIndex : "Source EIP was deleted before destination activation: operations=${migrationEipOperations}" } - applyEipOnDestHostSuccessed = false + migrationEipOperations.clear() env.afterSimulator(FlatEipBackend.BATCH_DELETE_EIP_PATH) { rsp, HttpEntity e -> - deleteEipOnSrcHostFailed = false + migrationEipOperations.add("delete") return rsp } @@ -279,8 +298,79 @@ class StartFlatNetworkVmWithEipCase extends SubCase { } retryInSecs { - assert deleteEipOnSrcHostFailed == false - assert applyEipOnDestHostSuccessed == true + int prepareIndex = migrationEipOperations.indexOf("prepare") + int activeIndex = migrationEipOperations.indexOf("active") + int deleteIndex = migrationEipOperations.indexOf("delete") + assert prepareIndex >= 0 : "Reverse migration did not prepare destination EIP: operations=${migrationEipOperations}" + assert activeIndex > prepareIndex : "Reverse migration did not activate destination after prepare: operations=${migrationEipOperations}" + assert deleteIndex > activeIndex : "Reverse migration deleted source before destination activation: operations=${migrationEipOperations}" + } + } + + void testAbnormalMigrationSourceCleanupFailure() { + def vm = env.inventoryByName("vm-1") as VmInstanceInventory + def host1 = env.inventoryByName("kvm") as HostInventory + def host2 = env.inventoryByName("kvm2") as HostInventory + + [ + VmAbnormalLifeCycleStruct.VmAbnormalLifeCycleOperation.VmMigrateToAnotherHost, + VmAbnormalLifeCycleStruct.VmAbnormalLifeCycleOperation.VmRunningFromUnknownStateHostChanged, + ].each { operation -> + AtomicInteger deleteCount = new AtomicInteger() + AtomicBoolean downstreamReached = new AtomicBoolean(false) + FutureCompletion chainResult = new FutureCompletion(null) + + env.afterSimulator(FlatEipBackend.BATCH_APPLY_EIP_PATH) { rsp -> + return rsp + } + env.afterSimulator(FlatEipBackend.BATCH_DELETE_EIP_PATH) { rsp -> + deleteCount.incrementAndGet() + rsp.success = false + rsp.error = "source cleanup failure" + return rsp + } + + VmAbnormalLifeCycleStruct struct = new VmAbnormalLifeCycleStruct() + struct.operation = operation + struct.vmInstance = org.zstack.header.vm.VmInstanceInventory.valueOf( + dbFindByUuid(vm.uuid, VmInstanceVO.class)) + struct.originalHostUuid = host1.uuid + struct.currentHostUuid = host2.uuid + + def chain = FlowChainBuilder.newSimpleFlowChain() + chain.then(bean(FlatEipBackend.class).createVmAbnormalLifeCycleHandlingFlow(struct)) + chain.then(new Flow() { + @Override + void run(FlowTrigger trigger, Map data) { + downstreamReached.set(true) + trigger.fail(operr("downstream failure")) + } + + @Override + void rollback(FlowRollback trigger, Map data) { + trigger.rollback() + } + }) + chain.done(new FlowDoneHandler(null) { + @Override + void handle(Map data) { + chainResult.success() + } + }) + chain.error(new FlowErrorHandler(null) { + @Override + void handle(org.zstack.header.errorcode.ErrorCode errorCode, Map data) { + chainResult.fail(errorCode) + } + }) + chain.start() + + chainResult.await(TimeUnit.SECONDS.toMillis(30)) + assert !chainResult.success + assert downstreamReached.get() : "source cleanup failure stopped ${operation} flow" + TimeUnit.MILLISECONDS.sleep(500) + assert deleteCount.get() == 1 : + "rollback deleted current-host EIP after source cleanup failure: operation=${operation}, deletes=${deleteCount.get()}" } }