Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -120,6 +122,7 @@ public static class DeleteEipCmd extends AgentCmd {

public static class BatchApplyEipCmd extends AgentCmd {
public List<EipTO> eips;
public boolean prepare;
}

public static class BatchDeleteEipCmd extends AgentCmd {
Expand All @@ -133,6 +136,17 @@ public static class BatchDeleteEipCmd extends AgentCmd {

@Override
public void preMigrateVm(VmInstanceInventory inv, String destHostUuid) {
List<EipTO> 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
Expand All @@ -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<EipTO> 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)
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -519,12 +565,18 @@ public void run(MessageReply reply) {
}

private void batchApplyEips(List<EipTO> eips, String hostUuid, final Completion completion) {
batchApplyEips(eips, hostUuid, false, completion);
batchApplyEips(eips, hostUuid, false, false, completion);
}

private void batchApplyEips(List<EipTO> eips, String hostUuid, boolean noHostStatusCheck, final Completion completion) {
batchApplyEips(eips, hostUuid, false, noHostStatusCheck, completion);
}

private void batchApplyEips(List<EipTO> 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);
Expand All @@ -547,6 +599,11 @@ public void run(MessageReply reply) {
return;
}

if (prepare) {
completion.success();
return;
}

List<String> vipUuids = CollectionUtils.transformToList(eips, new Function<String, EipTO>() {
@Override
public String call(EipTO arg) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

/**
Expand Down Expand Up @@ -245,15 +258,17 @@ class StartFlatNetworkVmWithEipCase extends SubCase {
vmNicUuid = vm.vmNics[0].uuid
}

boolean deleteEipOnSrcHostFailed = false
testAbnormalMigrationSourceCleanupFailure()

List<String> 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<String> e ->
applyEipOnDestHostSuccessed = true
def cmd = json(e.body, FlatEipBackend.BatchApplyEipCmd.class)
migrationEipOperations.add(cmd.prepare ? "prepare" : "active")
return rsp
}

Expand All @@ -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<String> e ->
deleteEipOnSrcHostFailed = false
migrationEipOperations.add("delete")
return rsp
}

Expand All @@ -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()}"
}
}

Expand Down