From 82225ade00d7706ddce77c4950cab64fe225c799 Mon Sep 17 00:00:00 2001 From: "shixin.ruan" Date: Thu, 20 Aug 2026 11:58:20 +0900 Subject: [PATCH 1/3] [flat-eip]: prepare EIP before migration Prepare destination EIPs before migration and let migrated events switch only the matching VM EIPs. Resolves: ZSTAC-87765 Change-Id: Ia7b2f4d1555851f6b73584327599776e27de4597 --- .../org/zstack/compute/vm/VmInstanceBase.java | 35 +++++--- .../vm/VmInstanceExtensionPointEmitter.java | 55 ++++++++++--- .../vm/VmMigrateCallExtensionFlow.java | 31 +++---- .../vm/VmInstanceMigrateExtensionPoint.java | 27 ++++++- .../network/service/flat/FlatEipBackend.java | 81 +++++++++++++------ 5 files changed, 159 insertions(+), 70 deletions(-) diff --git a/compute/src/main/java/org/zstack/compute/vm/VmInstanceBase.java b/compute/src/main/java/org/zstack/compute/vm/VmInstanceBase.java index 3f6328caef1..db29dff3023 100755 --- a/compute/src/main/java/org/zstack/compute/vm/VmInstanceBase.java +++ b/compute/src/main/java/org/zstack/compute/vm/VmInstanceBase.java @@ -6729,25 +6729,35 @@ public void done() { @Override public void handle(final Map data) { VmInstanceInventory vm = VmInstanceInventory.valueOf(self); - extEmitter.afterMigrateVm(vm, vm.getLastHostUuid()); - completion.success(); + extEmitter.afterMigrateVm(vm, vm.getLastHostUuid(), new NoErrorCompletion(completion) { + @Override + public void done() { + completion.success(); + } + }); } }).error(new FlowErrorHandler(completion) { @Override public void handle(final ErrorCode errCode, Map data) { String destHostUuid = spec.getDestHost().getUuid().equals(lastHostUuid) ? null : spec.getDestHost().getUuid(); - extEmitter.failedToMigrateVm(VmInstanceInventory.valueOf(self), destHostUuid, errCode); - if (HostErrors.FAILED_TO_MIGRATE_VM_ON_HYPERVISOR.isEqual(errCode.getCode())) { - checkState(originalCopy.getHostUuid(), new NoErrorCompletion(completion) { - @Override - public void done() { + extEmitter.failedToMigrateVm(VmInstanceInventory.valueOf(self), destHostUuid, errCode, + new NoErrorCompletion(completion) { + @Override + public void done() { + if (!HostErrors.FAILED_TO_MIGRATE_VM_ON_HYPERVISOR.isEqual(errCode.getCode())) { + changeVmStateInDb(originState.getDrivenEvent()); completion.fail(errCode); + return; } - }); - } else { - changeVmStateInDb(originState.getDrivenEvent()); - completion.fail(errCode); - } + + checkState(originalCopy.getHostUuid(), new NoErrorCompletion(completion) { + @Override + public void done() { + completion.fail(errCode); + } + }); + } + }); } }).start(); } @@ -8728,4 +8738,3 @@ public void run(MessageReply reply) { }); } } - diff --git a/compute/src/main/java/org/zstack/compute/vm/VmInstanceExtensionPointEmitter.java b/compute/src/main/java/org/zstack/compute/vm/VmInstanceExtensionPointEmitter.java index 354c303e003..ea25539420d 100755 --- a/compute/src/main/java/org/zstack/compute/vm/VmInstanceExtensionPointEmitter.java +++ b/compute/src/main/java/org/zstack/compute/vm/VmInstanceExtensionPointEmitter.java @@ -1,13 +1,17 @@ package org.zstack.compute.vm; import org.springframework.beans.factory.annotation.Autowired; +import org.zstack.core.asyncbatch.While; import org.zstack.core.componentloader.PluginRegistry; import org.zstack.core.errorcode.ErrorFacade; import org.zstack.core.workflow.FlowChainBuilder; import org.zstack.header.Component; import org.zstack.header.core.Completion; +import org.zstack.header.core.NoErrorCompletion; +import org.zstack.header.core.WhileDoneCompletion; import org.zstack.header.core.workflow.*; import org.zstack.header.errorcode.ErrorCode; +import org.zstack.header.errorcode.ErrorCodeList; import org.zstack.header.errorcode.SysErrors; import org.zstack.header.vm.*; import org.zstack.header.volume.VolumeInventory; @@ -320,28 +324,59 @@ public void run(VmInstanceStartExtensionPoint arg) { }); } - public void preMigrateVm(final VmInstanceInventory inv, final String dstHostUuid) { - CollectionUtils.safeForEach(migrateVmExtensions, arg -> arg.preMigrateVm(inv, dstHostUuid)); + public void preMigrateVm(final VmInstanceInventory inv, final String dstHostUuid, Completion completion) { + new While<>(migrateVmExtensions).each((ext, comp) -> ext.preMigrateVm(inv, dstHostUuid, new Completion(comp) { + @Override + public void success() { + comp.done(); + } + + @Override + public void fail(ErrorCode errorCode) { + comp.addError(errorCode); + comp.allDone(); + } + })).run(new WhileDoneCompletion(completion) { + @Override + public void done(ErrorCodeList errorCodeList) { + if (errorCodeList.getCauses().isEmpty()) { + completion.success(); + } else { + completion.fail(errorCodeList.getCauses().get(0)); + } + } + }); } public void beforeMigrateVm(final VmInstanceInventory inv, final String dstHostUuid) { CollectionUtils.safeForEach(migrateVmExtensions, arg -> arg.beforeMigrateVm(inv, dstHostUuid)); } - public void afterMigrateVm(final VmInstanceInventory inv, final String srcHostUuid) { - CollectionUtils.safeForEach(migrateVmExtensions, new ForEachFunction() { + public void afterMigrateVm(final VmInstanceInventory inv, final String srcHostUuid, NoErrorCompletion completion) { + new While<>(migrateVmExtensions).each((ext, comp) -> ext.afterMigrateVm(inv, srcHostUuid, new NoErrorCompletion(comp) { + @Override + public void done() { + comp.done(); + } + })).run(new WhileDoneCompletion(completion) { @Override - public void run(VmInstanceMigrateExtensionPoint arg) { - arg.afterMigrateVm(inv, srcHostUuid); + public void done(ErrorCodeList errorCodeList) { + completion.done(); } }); } - public void failedToMigrateVm(final VmInstanceInventory inv, final String dstHostUuid, final ErrorCode reason) { - CollectionUtils.safeForEach(migrateVmExtensions, new ForEachFunction() { + public void failedToMigrateVm(final VmInstanceInventory inv, final String dstHostUuid, final ErrorCode reason, + NoErrorCompletion completion) { + new While<>(migrateVmExtensions).each((ext, comp) -> ext.failedToMigrateVm(inv, dstHostUuid, reason, new NoErrorCompletion(comp) { + @Override + public void done() { + comp.done(); + } + })).run(new WhileDoneCompletion(completion) { @Override - public void run(final VmInstanceMigrateExtensionPoint arg) { - arg.failedToMigrateVm(inv, dstHostUuid, reason); + public void done(ErrorCodeList errorCodeList) { + completion.done(); } }); } diff --git a/compute/src/main/java/org/zstack/compute/vm/VmMigrateCallExtensionFlow.java b/compute/src/main/java/org/zstack/compute/vm/VmMigrateCallExtensionFlow.java index d2901d18341..9349a99a45a 100755 --- a/compute/src/main/java/org/zstack/compute/vm/VmMigrateCallExtensionFlow.java +++ b/compute/src/main/java/org/zstack/compute/vm/VmMigrateCallExtensionFlow.java @@ -3,17 +3,13 @@ import org.springframework.beans.factory.annotation.Autowire; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Configurable; -import org.zstack.core.componentloader.PluginRegistry; +import org.zstack.header.core.Completion; import org.zstack.header.core.workflow.Flow; import org.zstack.header.core.workflow.FlowRollback; import org.zstack.header.core.workflow.FlowTrigger; import org.zstack.header.host.HostInventory; -import org.zstack.header.vm.MigrateVmMessage; import org.zstack.header.vm.VmInstanceConstant; -import org.zstack.header.vm.VmInstanceMigrateExtensionPoint; import org.zstack.header.vm.VmInstanceSpec; -import org.zstack.utils.CollectionUtils; -import org.zstack.utils.function.ForEachFunction; import java.util.Map; @@ -22,32 +18,25 @@ @Configurable(preConstruction = true, autowire = Autowire.BY_TYPE) public class VmMigrateCallExtensionFlow implements Flow { @Autowired - protected PluginRegistry pluginRgty; + private VmInstanceExtensionPointEmitter extEmitter; @Override public void run(FlowTrigger trigger, Map data) { final VmInstanceSpec spec = (VmInstanceSpec) data.get(VmInstanceConstant.Params.VmInstanceSpec.toString()); - boolean migrateFromDest = false; - if (spec.getMessage() instanceof MigrateVmMessage) { - migrateFromDest = ((MigrateVmMessage)spec.getMessage()).isMigrateFromDestination(); - } - - final boolean fromDest = migrateFromDest; - final HostInventory destHost = spec.getDestHost(); - for (VmInstanceMigrateExtensionPoint ext : pluginRgty.getExtensionList(VmInstanceMigrateExtensionPoint.class)) { - ext.preMigrateVm(spec.getVmInventory(), destHost.getUuid()); - } + extEmitter.preMigrateVm(spec.getVmInventory(), destHost.getUuid(), new Completion(trigger) { + @Override + public void success() { + extEmitter.beforeMigrateVm(spec.getVmInventory(), destHost.getUuid()); + trigger.next(); + } - CollectionUtils.safeForEach(pluginRgty.getExtensionList(VmInstanceMigrateExtensionPoint.class), new ForEachFunction() { @Override - public void run(VmInstanceMigrateExtensionPoint ext) { - ext.beforeMigrateVm(spec.getVmInventory(), destHost.getUuid()); + public void fail(org.zstack.header.errorcode.ErrorCode errorCode) { + trigger.fail(errorCode); } }); - - trigger.next(); } @Override diff --git a/header/src/main/java/org/zstack/header/vm/VmInstanceMigrateExtensionPoint.java b/header/src/main/java/org/zstack/header/vm/VmInstanceMigrateExtensionPoint.java index a5576c9aa92..32fe130d126 100755 --- a/header/src/main/java/org/zstack/header/vm/VmInstanceMigrateExtensionPoint.java +++ b/header/src/main/java/org/zstack/header/vm/VmInstanceMigrateExtensionPoint.java @@ -1,13 +1,34 @@ package org.zstack.header.vm; +import org.zstack.header.core.Completion; +import org.zstack.header.core.NoErrorCompletion; import org.zstack.header.errorcode.ErrorCode; public interface VmInstanceMigrateExtensionPoint { - void preMigrateVm(VmInstanceInventory inv, String destHostUuid); + default void preMigrateVm(VmInstanceInventory inv, String destHostUuid) { + } + + default void preMigrateVm(VmInstanceInventory inv, String destHostUuid, Completion completion) { + preMigrateVm(inv, destHostUuid); + completion.success(); + } void beforeMigrateVm(VmInstanceInventory inv, String destHostUuid); - void afterMigrateVm(VmInstanceInventory inv, String srcHostUuid); + default void afterMigrateVm(VmInstanceInventory inv, String srcHostUuid) { + } + + default void afterMigrateVm(VmInstanceInventory inv, String srcHostUuid, NoErrorCompletion completion) { + afterMigrateVm(inv, srcHostUuid); + completion.done(); + } + + default void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason) { + } - void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason); + default void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason, + NoErrorCompletion completion) { + failedToMigrateVm(inv, destHostUuid, reason); + completion.done(); + } } 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..15c18bd3476 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.NoErrorCompletion; import org.zstack.header.core.NopeCompletion; import org.zstack.header.core.workflow.Flow; import org.zstack.header.core.workflow.FlowRollback; @@ -120,6 +121,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 { @@ -132,7 +134,14 @@ public static class BatchDeleteEipCmd extends AgentCmd { public static final String BATCH_DELETE_EIP_PATH = "/flatnetworkprovider/eip/batchdelete"; @Override - public void preMigrateVm(VmInstanceInventory inv, String destHostUuid) { + public void preMigrateVm(VmInstanceInventory inv, String destHostUuid, Completion completion) { + List eips = getEipsByVmUuid(inv.getUuid()); + if (eips == null || eips.isEmpty()) { + completion.success(); + return; + } + + batchApplyEips(eips, destHostUuid, true, false, completion); } @Override @@ -141,52 +150,67 @@ public void beforeMigrateVm(VmInstanceInventory inv, String destHostUuid) { } @Override - public void afterMigrateVm(final VmInstanceInventory inv, String srcHostUuid) { + public void afterMigrateVm(final VmInstanceInventory inv, String srcHostUuid, NoErrorCompletion completion) { List eips = getEipsByVmUuid(inv.getUuid()); if (eips == null || eips.isEmpty()) { + completion.done(); return; } - batchDeleteEips(eips, srcHostUuid, new Completion(null) { + batchApplyEips(eips, inv.getHostUuid(), new Completion(completion) { @Override public void success() { - batchApplyEips(eips, inv.getHostUuid(), new Completion(null) { + batchDeleteEips(eips, srcHostUuid, new Completion(completion) { @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)); + completion.done(); } @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)); + completion.done(); } }); } @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)); + completion.done(); } }); } @Override - public void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason) { + public void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason, + NoErrorCompletion completion) { + List eips = getEipsByVmUuid(inv.getUuid()); + if (eips == null || eips.isEmpty() || destHostUuid == null) { + completion.done(); + return; + } + if (destHostUuid.equals(inv.getHostUuid())) { + afterMigrateVm(inv, inv.getLastHostUuid(), completion); + return; + } + + batchDeleteEips(eips, destHostUuid, new Completion(completion) { + @Override + public void success() { + completion.done(); + } + + @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)); + completion.done(); + } + }); } @Transactional(readOnly = true) @@ -519,12 +543,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 +577,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) { From 1fa3c2725d225d5dd48afbb0d2fcee6c1aca96b2 Mon Sep 17 00:00:00 2001 From: "shixin.ruan" Date: Thu, 20 Aug 2026 20:58:35 +0900 Subject: [PATCH 2/3] [flat-eip]: finalize migration handoff Keep legacy migration callbacks unchanged.\nPrepare destination EIP before migration and activate it before source cleanup. Change-Id: Iabf21c09ff9f95c3866b03fe91fefe9b410db82e --- .../org/zstack/compute/vm/VmInstanceBase.java | 35 ++++------ .../vm/VmInstanceExtensionPointEmitter.java | 55 +++------------- .../vm/VmMigrateCallExtensionFlow.java | 31 ++++++--- .../vm/VmInstanceMigrateExtensionPoint.java | 27 +------- .../network/service/flat/FlatEipBackend.java | 64 ++++++++++++------- .../eip/StartFlatNetworkVmWithEipCase.groovy | 28 +++++--- 6 files changed, 105 insertions(+), 135 deletions(-) diff --git a/compute/src/main/java/org/zstack/compute/vm/VmInstanceBase.java b/compute/src/main/java/org/zstack/compute/vm/VmInstanceBase.java index db29dff3023..3f6328caef1 100755 --- a/compute/src/main/java/org/zstack/compute/vm/VmInstanceBase.java +++ b/compute/src/main/java/org/zstack/compute/vm/VmInstanceBase.java @@ -6729,35 +6729,25 @@ public void done() { @Override public void handle(final Map data) { VmInstanceInventory vm = VmInstanceInventory.valueOf(self); - extEmitter.afterMigrateVm(vm, vm.getLastHostUuid(), new NoErrorCompletion(completion) { - @Override - public void done() { - completion.success(); - } - }); + extEmitter.afterMigrateVm(vm, vm.getLastHostUuid()); + completion.success(); } }).error(new FlowErrorHandler(completion) { @Override public void handle(final ErrorCode errCode, Map data) { String destHostUuid = spec.getDestHost().getUuid().equals(lastHostUuid) ? null : spec.getDestHost().getUuid(); - extEmitter.failedToMigrateVm(VmInstanceInventory.valueOf(self), destHostUuid, errCode, - new NoErrorCompletion(completion) { - @Override - public void done() { - if (!HostErrors.FAILED_TO_MIGRATE_VM_ON_HYPERVISOR.isEqual(errCode.getCode())) { - changeVmStateInDb(originState.getDrivenEvent()); + extEmitter.failedToMigrateVm(VmInstanceInventory.valueOf(self), destHostUuid, errCode); + if (HostErrors.FAILED_TO_MIGRATE_VM_ON_HYPERVISOR.isEqual(errCode.getCode())) { + checkState(originalCopy.getHostUuid(), new NoErrorCompletion(completion) { + @Override + public void done() { completion.fail(errCode); - return; } - - checkState(originalCopy.getHostUuid(), new NoErrorCompletion(completion) { - @Override - public void done() { - completion.fail(errCode); - } - }); - } - }); + }); + } else { + changeVmStateInDb(originState.getDrivenEvent()); + completion.fail(errCode); + } } }).start(); } @@ -8738,3 +8728,4 @@ public void run(MessageReply reply) { }); } } + diff --git a/compute/src/main/java/org/zstack/compute/vm/VmInstanceExtensionPointEmitter.java b/compute/src/main/java/org/zstack/compute/vm/VmInstanceExtensionPointEmitter.java index ea25539420d..354c303e003 100755 --- a/compute/src/main/java/org/zstack/compute/vm/VmInstanceExtensionPointEmitter.java +++ b/compute/src/main/java/org/zstack/compute/vm/VmInstanceExtensionPointEmitter.java @@ -1,17 +1,13 @@ package org.zstack.compute.vm; import org.springframework.beans.factory.annotation.Autowired; -import org.zstack.core.asyncbatch.While; import org.zstack.core.componentloader.PluginRegistry; import org.zstack.core.errorcode.ErrorFacade; import org.zstack.core.workflow.FlowChainBuilder; import org.zstack.header.Component; import org.zstack.header.core.Completion; -import org.zstack.header.core.NoErrorCompletion; -import org.zstack.header.core.WhileDoneCompletion; import org.zstack.header.core.workflow.*; import org.zstack.header.errorcode.ErrorCode; -import org.zstack.header.errorcode.ErrorCodeList; import org.zstack.header.errorcode.SysErrors; import org.zstack.header.vm.*; import org.zstack.header.volume.VolumeInventory; @@ -324,59 +320,28 @@ public void run(VmInstanceStartExtensionPoint arg) { }); } - public void preMigrateVm(final VmInstanceInventory inv, final String dstHostUuid, Completion completion) { - new While<>(migrateVmExtensions).each((ext, comp) -> ext.preMigrateVm(inv, dstHostUuid, new Completion(comp) { - @Override - public void success() { - comp.done(); - } - - @Override - public void fail(ErrorCode errorCode) { - comp.addError(errorCode); - comp.allDone(); - } - })).run(new WhileDoneCompletion(completion) { - @Override - public void done(ErrorCodeList errorCodeList) { - if (errorCodeList.getCauses().isEmpty()) { - completion.success(); - } else { - completion.fail(errorCodeList.getCauses().get(0)); - } - } - }); + public void preMigrateVm(final VmInstanceInventory inv, final String dstHostUuid) { + CollectionUtils.safeForEach(migrateVmExtensions, arg -> arg.preMigrateVm(inv, dstHostUuid)); } public void beforeMigrateVm(final VmInstanceInventory inv, final String dstHostUuid) { CollectionUtils.safeForEach(migrateVmExtensions, arg -> arg.beforeMigrateVm(inv, dstHostUuid)); } - public void afterMigrateVm(final VmInstanceInventory inv, final String srcHostUuid, NoErrorCompletion completion) { - new While<>(migrateVmExtensions).each((ext, comp) -> ext.afterMigrateVm(inv, srcHostUuid, new NoErrorCompletion(comp) { - @Override - public void done() { - comp.done(); - } - })).run(new WhileDoneCompletion(completion) { + public void afterMigrateVm(final VmInstanceInventory inv, final String srcHostUuid) { + CollectionUtils.safeForEach(migrateVmExtensions, new ForEachFunction() { @Override - public void done(ErrorCodeList errorCodeList) { - completion.done(); + public void run(VmInstanceMigrateExtensionPoint arg) { + arg.afterMigrateVm(inv, srcHostUuid); } }); } - public void failedToMigrateVm(final VmInstanceInventory inv, final String dstHostUuid, final ErrorCode reason, - NoErrorCompletion completion) { - new While<>(migrateVmExtensions).each((ext, comp) -> ext.failedToMigrateVm(inv, dstHostUuid, reason, new NoErrorCompletion(comp) { - @Override - public void done() { - comp.done(); - } - })).run(new WhileDoneCompletion(completion) { + public void failedToMigrateVm(final VmInstanceInventory inv, final String dstHostUuid, final ErrorCode reason) { + CollectionUtils.safeForEach(migrateVmExtensions, new ForEachFunction() { @Override - public void done(ErrorCodeList errorCodeList) { - completion.done(); + public void run(final VmInstanceMigrateExtensionPoint arg) { + arg.failedToMigrateVm(inv, dstHostUuid, reason); } }); } diff --git a/compute/src/main/java/org/zstack/compute/vm/VmMigrateCallExtensionFlow.java b/compute/src/main/java/org/zstack/compute/vm/VmMigrateCallExtensionFlow.java index 9349a99a45a..d2901d18341 100755 --- a/compute/src/main/java/org/zstack/compute/vm/VmMigrateCallExtensionFlow.java +++ b/compute/src/main/java/org/zstack/compute/vm/VmMigrateCallExtensionFlow.java @@ -3,13 +3,17 @@ import org.springframework.beans.factory.annotation.Autowire; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Configurable; -import org.zstack.header.core.Completion; +import org.zstack.core.componentloader.PluginRegistry; import org.zstack.header.core.workflow.Flow; import org.zstack.header.core.workflow.FlowRollback; import org.zstack.header.core.workflow.FlowTrigger; import org.zstack.header.host.HostInventory; +import org.zstack.header.vm.MigrateVmMessage; import org.zstack.header.vm.VmInstanceConstant; +import org.zstack.header.vm.VmInstanceMigrateExtensionPoint; import org.zstack.header.vm.VmInstanceSpec; +import org.zstack.utils.CollectionUtils; +import org.zstack.utils.function.ForEachFunction; import java.util.Map; @@ -18,25 +22,32 @@ @Configurable(preConstruction = true, autowire = Autowire.BY_TYPE) public class VmMigrateCallExtensionFlow implements Flow { @Autowired - private VmInstanceExtensionPointEmitter extEmitter; + protected PluginRegistry pluginRgty; @Override public void run(FlowTrigger trigger, Map data) { final VmInstanceSpec spec = (VmInstanceSpec) data.get(VmInstanceConstant.Params.VmInstanceSpec.toString()); + boolean migrateFromDest = false; + if (spec.getMessage() instanceof MigrateVmMessage) { + migrateFromDest = ((MigrateVmMessage)spec.getMessage()).isMigrateFromDestination(); + } + + final boolean fromDest = migrateFromDest; + final HostInventory destHost = spec.getDestHost(); - extEmitter.preMigrateVm(spec.getVmInventory(), destHost.getUuid(), new Completion(trigger) { - @Override - public void success() { - extEmitter.beforeMigrateVm(spec.getVmInventory(), destHost.getUuid()); - trigger.next(); - } + for (VmInstanceMigrateExtensionPoint ext : pluginRgty.getExtensionList(VmInstanceMigrateExtensionPoint.class)) { + ext.preMigrateVm(spec.getVmInventory(), destHost.getUuid()); + } + CollectionUtils.safeForEach(pluginRgty.getExtensionList(VmInstanceMigrateExtensionPoint.class), new ForEachFunction() { @Override - public void fail(org.zstack.header.errorcode.ErrorCode errorCode) { - trigger.fail(errorCode); + public void run(VmInstanceMigrateExtensionPoint ext) { + ext.beforeMigrateVm(spec.getVmInventory(), destHost.getUuid()); } }); + + trigger.next(); } @Override diff --git a/header/src/main/java/org/zstack/header/vm/VmInstanceMigrateExtensionPoint.java b/header/src/main/java/org/zstack/header/vm/VmInstanceMigrateExtensionPoint.java index 32fe130d126..a5576c9aa92 100755 --- a/header/src/main/java/org/zstack/header/vm/VmInstanceMigrateExtensionPoint.java +++ b/header/src/main/java/org/zstack/header/vm/VmInstanceMigrateExtensionPoint.java @@ -1,34 +1,13 @@ package org.zstack.header.vm; -import org.zstack.header.core.Completion; -import org.zstack.header.core.NoErrorCompletion; import org.zstack.header.errorcode.ErrorCode; public interface VmInstanceMigrateExtensionPoint { - default void preMigrateVm(VmInstanceInventory inv, String destHostUuid) { - } - - default void preMigrateVm(VmInstanceInventory inv, String destHostUuid, Completion completion) { - preMigrateVm(inv, destHostUuid); - completion.success(); - } + void preMigrateVm(VmInstanceInventory inv, String destHostUuid); void beforeMigrateVm(VmInstanceInventory inv, String destHostUuid); - default void afterMigrateVm(VmInstanceInventory inv, String srcHostUuid) { - } - - default void afterMigrateVm(VmInstanceInventory inv, String srcHostUuid, NoErrorCompletion completion) { - afterMigrateVm(inv, srcHostUuid); - completion.done(); - } - - default void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason) { - } + void afterMigrateVm(VmInstanceInventory inv, String srcHostUuid); - default void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason, - NoErrorCompletion completion) { - failedToMigrateVm(inv, destHostUuid, reason); - completion.done(); - } + void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason); } 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 15c18bd3476..fa0d2fccdbb 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,7 +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.NoErrorCompletion; +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; @@ -58,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; @@ -134,14 +135,18 @@ public static class BatchDeleteEipCmd extends AgentCmd { public static final String BATCH_DELETE_EIP_PATH = "/flatnetworkprovider/eip/batchdelete"; @Override - public void preMigrateVm(VmInstanceInventory inv, String destHostUuid, Completion completion) { + public void preMigrateVm(VmInstanceInventory inv, String destHostUuid) { List eips = getEipsByVmUuid(inv.getUuid()); if (eips == null || eips.isEmpty()) { - completion.success(); 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 @@ -150,27 +155,24 @@ public void beforeMigrateVm(VmInstanceInventory inv, String destHostUuid) { } @Override - public void afterMigrateVm(final VmInstanceInventory inv, String srcHostUuid, NoErrorCompletion completion) { + public void afterMigrateVm(final VmInstanceInventory inv, String srcHostUuid) { List eips = getEipsByVmUuid(inv.getUuid()); if (eips == null || eips.isEmpty()) { - completion.done(); return; } - batchApplyEips(eips, inv.getHostUuid(), new Completion(completion) { + batchApplyEips(eips, inv.getHostUuid(), new Completion(null) { @Override public void success() { - batchDeleteEips(eips, srcHostUuid, new Completion(completion) { + batchDeleteEips(eips, srcHostUuid, new Completion(null) { @Override public void success() { - completion.done(); } @Override public void fail(ErrorCode 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)); - completion.done(); } }); } @@ -179,36 +181,32 @@ public void fail(ErrorCode errorCode) { public void fail(ErrorCode 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)); - completion.done(); } }); } @Override - public void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason, - NoErrorCompletion completion) { + public void failedToMigrateVm(VmInstanceInventory inv, String destHostUuid, ErrorCode reason) { List eips = getEipsByVmUuid(inv.getUuid()); if (eips == null || eips.isEmpty() || destHostUuid == null) { - completion.done(); return; } - if (destHostUuid.equals(inv.getHostUuid())) { - afterMigrateVm(inv, inv.getLastHostUuid(), completion); + 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(completion) { + batchDeleteEips(eips, destHostUuid, new Completion(null) { @Override public void success() { - completion.done(); } @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)); - completion.done(); } }); } @@ -274,13 +272,22 @@ 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) { + trigger.fail(errorCode); + } + }); } @Override @@ -292,13 +299,22 @@ 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) { + trigger.fail(errorCode); + } + }); } @Override 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..23e07e0972c 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 @@ -245,15 +245,15 @@ class StartFlatNetworkVmWithEipCase extends SubCase { vmNicUuid = vm.vmNics[0].uuid } - boolean deleteEipOnSrcHostFailed = false + 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 +263,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 +283,12 @@ 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}" } } From f02510bd2ded08d1670b45d9becc0a394330dfb1 Mon Sep 17 00:00:00 2001 From: "shixin.ruan" Date: Thu, 20 Aug 2026 23:10:34 +0900 Subject: [PATCH 3/3] [flat-eip]: preserve EIP after cleanup failure Keep the current-host EIP when original-host cleanup fails after an abnormal migration. Resolves: ZSTAC-87765 Change-Id: I129d25624275e08b6d5f6bb4f34e27dc47c7809c --- .../network/service/flat/FlatEipBackend.java | 10 ++- .../eip/StartFlatNetworkVmWithEipCase.groovy | 82 +++++++++++++++++++ 2 files changed, 90 insertions(+), 2 deletions(-) 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 fa0d2fccdbb..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 @@ -285,7 +285,10 @@ public void success() { @Override public void fail(ErrorCode errorCode) { - trigger.fail(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(); } }); } @@ -312,7 +315,10 @@ public void success() { @Override public void fail(ErrorCode errorCode) { - trigger.fail(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(); } }); } 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 23e07e0972c..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,6 +258,8 @@ class StartFlatNetworkVmWithEipCase extends SubCase { vmNicUuid = vm.vmNics[0].uuid } + testAbnormalMigrationSourceCleanupFailure() + List migrationEipOperations = Collections.synchronizedList([]) env.afterSimulator(FlatEipBackend.BATCH_DELETE_EIP_PATH) { migrationEipOperations.add("delete") @@ -292,6 +307,73 @@ class StartFlatNetworkVmWithEipCase extends SubCase { } } + 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()}" + } + } + @Override void clean() { env.delete()