From a49148c7e4dadb06c2cebe6111ecedb540662578 Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Tue, 19 Sep 2023 16:29:34 +0800 Subject: [PATCH] fix(region,host): cancel migrate sync guest status --- pkg/compute/guestdrivers/base.go | 4 ++++ pkg/compute/guestdrivers/kvm.go | 12 +++++++++++ pkg/compute/models/guest_actions.go | 7 ++----- pkg/compute/models/guestdrivers.go | 1 + pkg/compute/tasks/guest_live_migrate_task.go | 19 +++++++++++++++++- .../guestman/guesthandlers/guesthandler.go | 20 +++++++++++++++++++ pkg/hostman/guestman/guesttasks.go | 15 +++++++++++--- pkg/hostman/monitor/hmp.go | 4 ++++ pkg/hostman/monitor/monitor.go | 1 + pkg/hostman/monitor/qmp.go | 4 ++++ 10 files changed, 78 insertions(+), 9 deletions(-) diff --git a/pkg/compute/guestdrivers/base.go b/pkg/compute/guestdrivers/base.go index b4d1bbdc1b..772f6e2104 100644 --- a/pkg/compute/guestdrivers/base.go +++ b/pkg/compute/guestdrivers/base.go @@ -441,6 +441,10 @@ func (drv *SBaseGuestDriver) RequestLiveMigrate(ctx context.Context, guest *mode return errors.Wrapf(cloudprovider.ErrNotImplemented, "RequestLiveMigrate") } +func (drv *SBaseGuestDriver) RequestCancelLiveMigrate(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential) error { + return errors.Wrapf(cloudprovider.ErrNotImplemented, "RequestCancelLiveMigrate") +} + func (drv *SVirtualizedGuestDriver) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, input *api.ServerCreateInput) (*api.ServerCreateInput, error) { return input, nil } diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index 50374b88bc..dcd4602933 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -845,6 +845,18 @@ func (self *SKVMGuestDriver) CheckLiveMigrate(ctx context.Context, guest *models return nil } +func (self *SKVMGuestDriver) RequestCancelLiveMigrate(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential) error { + host, _ := guest.GetHost() + url := fmt.Sprintf("%s/servers/%s/cancel-live-migrate", host.ManagerUri, guest.Id) + httpClient := httputils.GetDefaultClient() + header := mcclient.GetTokenHeaders(userCred) + _, _, err := httputils.JSONRequest(httpClient, ctx, "POST", url, header, jsonutils.NewDict(), false) + if err != nil { + return errors.Wrap(err, "host request") + } + return nil +} + func (self *SKVMGuestDriver) ValidateDetachNetwork(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest) error { if guest.Status == api.VM_RUNNING && guest.GetMetadata(ctx, "hot_remove_nic", nil) != "enable" { return httperrors.NewBadRequestError("Guest %s can't hot remove nic", guest.GetName()) diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 6534bbcbcb..da99d90af3 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -647,11 +647,8 @@ func (self *SGuest) PerformCancelLiveMigrate( if self.Status != api.VM_LIVE_MIGRATING { return nil, httperrors.NewServerStatusError("cannot set migrate params in status %s", self.Status) } - monitorInput := &api.ServerMonitorInput{ - COMMAND: "migrate_cancel", - QMP: false, - } - return self.SendMonitorCommand(ctx, userCred, monitorInput) + + return nil, self.GetDriver().RequestCancelLiveMigrate(ctx, self, userCred) } func (self *SGuest) PerformClone(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { diff --git a/pkg/compute/models/guestdrivers.go b/pkg/compute/models/guestdrivers.go index 4cb17f0c9b..405d177e4b 100644 --- a/pkg/compute/models/guestdrivers.go +++ b/pkg/compute/models/guestdrivers.go @@ -213,6 +213,7 @@ type IGuestDriver interface { CheckLiveMigrate(ctx context.Context, guest *SGuest, userCred mcclient.TokenCredential, input api.GuestLiveMigrateInput) error RequestMigrate(ctx context.Context, guest *SGuest, userCred mcclient.TokenCredential, input api.GuestMigrateInput, task taskman.ITask) error RequestLiveMigrate(ctx context.Context, guest *SGuest, userCred mcclient.TokenCredential, input api.GuestLiveMigrateInput, task taskman.ITask) error + RequestCancelLiveMigrate(ctx context.Context, guest *SGuest, userCred mcclient.TokenCredential) error ValidateUpdateData(ctx context.Context, guest *SGuest, userCred mcclient.TokenCredential, input api.ServerUpdateInput) (api.ServerUpdateInput, error) RequestRemoteUpdate(ctx context.Context, guest *SGuest, userCred mcclient.TokenCredential, replaceTags bool) error diff --git a/pkg/compute/tasks/guest_live_migrate_task.go b/pkg/compute/tasks/guest_live_migrate_task.go index 20e0b82677..774f033635 100644 --- a/pkg/compute/tasks/guest_live_migrate_task.go +++ b/pkg/compute/tasks/guest_live_migrate_task.go @@ -547,10 +547,27 @@ func (task *GuestMigrateTask) setGuest(ctx context.Context, guest *models.SGuest } func (task *GuestLiveMigrateTask) OnLiveMigrateCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + if reason, _ := data.GetString("__reason__"); reason == "cancelled" { + task.Params.Set("migrate_cancelled", jsonutils.JSONTrue) + } + if !jsonutils.QueryBoolean(task.Params, "keep_dest_guest_on_failed", false) { targetHostId, _ := task.Params.GetString("target_host_id") - guest.StartUndeployGuestTask(ctx, task.UserCred, "", targetHostId) + task.SetStage("OnGuestUndeployed", nil) + guest.StartUndeployGuestTask(ctx, task.UserCred, task.Id, targetHostId) + } else { + task.OnGuestUndeployed(ctx, guest, data) } +} + +func (task *GuestLiveMigrateTask) OnGuestUndeployed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + task.TaskFailed(ctx, guest, data) + if jsonutils.QueryBoolean(task.Params, "migrate_cancelled", false) { + guest.StartSyncstatus(ctx, task.UserCred, "") + } +} + +func (task *GuestLiveMigrateTask) OnGuestUndeployedFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { task.TaskFailed(ctx, guest, data) } diff --git a/pkg/hostman/guestman/guesthandlers/guesthandler.go b/pkg/hostman/guestman/guesthandlers/guesthandler.go index d36226f436..9ae687ff4d 100644 --- a/pkg/hostman/guestman/guesthandlers/guesthandler.go +++ b/pkg/hostman/guestman/guesthandlers/guesthandler.go @@ -77,6 +77,7 @@ func AddGuestTaskHandler(prefix string, app *appsrv.Application) { "src-prepare-migrate": guestSrcPrepareMigrate, "dest-prepare-migrate": guestDestPrepareMigrate, "live-migrate": guestLiveMigrate, + "cancel-live-migrate": guestCancelLiveMigrate, "resume": guestResume, "block-replication": guestBlockReplication, "slave-block-stream-disks": slaveGuestBlockStreamDisks, @@ -474,6 +475,25 @@ func guestLiveMigrate(ctx context.Context, userCred mcclient.TokenCredential, si return nil, nil } +func guestCancelLiveMigrate(ctx context.Context, userCred mcclient.TokenCredential, sid string, body jsonutils.JSONObject) (interface{}, error) { + guest, ok := guestman.GetGuestManager().GetServer(sid) + if !ok { + return nil, httperrors.NewNotFoundError("Guest %s not found", sid) + } + if guest.MigrateTask == nil { + return nil, httperrors.NewBadRequestError("Guest %s not in migrating", sid) + } + guest.MigrateTask.SetLiveMigrateCancelled() + var c = make(chan string) + cb := func(res string) { + c <- res + } + guest.Monitor.MigrateCancel(cb) + var res = <-c + lines := strings.Split(res, "\\r\\n") + return strDict{"results": strings.Join(lines, "\n")}, nil +} + func guestResume(ctx context.Context, userCred mcclient.TokenCredential, sid string, body jsonutils.JSONObject) (interface{}, error) { if !guestman.GetGuestManager().IsGuestExist(sid) { return nil, httperrors.NewNotFoundError("Guest %s not found", sid) diff --git a/pkg/hostman/guestman/guesttasks.go b/pkg/hostman/guestman/guesttasks.go index 5b2b0837b9..7f78810872 100644 --- a/pkg/hostman/guestman/guesttasks.go +++ b/pkg/hostman/guestman/guesttasks.go @@ -1003,6 +1003,7 @@ type SGuestLiveMigrateTask struct { timeoutAt time.Time doTimeoutMigrate bool + cancelled bool expectDowntime int64 dirtySyncCount int64 @@ -1255,12 +1256,14 @@ func (s *SGuestLiveMigrateTask) onGetMigrateStats(stats *monitor.MigrationInfo, */ func (s *SGuestLiveMigrateTask) onGetMigrateStatus(stats *monitor.MigrationInfo) { - status := *stats.Status + status := string(*stats.Status) if status == "completed" { jsonStats := jsonutils.Marshal(stats) log.Infof("migration info %s", jsonStats) - } else if status == "failed" || status == "cancelled" { + } else if status == "failed" { s.migrateFailed(fmt.Sprintf("Query migrate got status: %s", status)) + } else if status == "cancelled" { + s.migrateFailed(status) } else if status == "active" { var ( ramRemain int64 @@ -1335,7 +1338,9 @@ func (s *SGuestLiveMigrateTask) onMigrateReceivedPreSwitchoverEvent() { } func (s *SGuestLiveMigrateTask) onMigrateReceivedBlockJobError(res string) { - s.migrateFailed(res) + if !s.cancelled { + s.migrateFailed(res) + } } func (s *SGuestLiveMigrateTask) migrateComplete(stats jsonutils.JSONObject) { @@ -1396,6 +1401,10 @@ func (s *SGuestLiveMigrateTask) onMigrateFailBlockJobsCancelled(msg string) { } } +func (s *SGuestLiveMigrateTask) SetLiveMigrateCancelled() { + s.cancelled = true +} + /** * GuestResumeTask **/ diff --git a/pkg/hostman/monitor/hmp.go b/pkg/hostman/monitor/hmp.go index 8d9a1d3355..9f8f3194e1 100644 --- a/pkg/hostman/monitor/hmp.go +++ b/pkg/hostman/monitor/hmp.go @@ -359,6 +359,10 @@ func (m *HmpMonitor) GetMigrateStats(callback MigrateStatsCallback) { go callback(nil, errors.Errorf("unsupport get migrate stats")) } +func (m *HmpMonitor) MigrateCancel(cb StringCallback) { + m.Query("migrate_cancel", cb) +} + func (m *HmpMonitor) MigrateStartPostcopy(callback StringCallback) { cb := func(output string) { log.Infof("MigrateStartPostcopy %s: %s", m.server, output) diff --git a/pkg/hostman/monitor/monitor.go b/pkg/hostman/monitor/monitor.go index 792a947262..aab120663e 100644 --- a/pkg/hostman/monitor/monitor.go +++ b/pkg/hostman/monitor/monitor.go @@ -238,6 +238,7 @@ type Monitor interface { GetMigrateStatus(callback StringCallback) MigrateStartPostcopy(callback StringCallback) GetMigrateStats(callback MigrateStatsCallback) + MigrateCancel(cb StringCallback) ReloadDiskBlkdev(device, path string, callback StringCallback) SetVncPassword(proto, password string, callback StringCallback) diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index 0e22c488b8..99ba909746 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -687,6 +687,10 @@ func (m *QmpMonitor) GetMigrateStats(callback MigrateStatsCallback) { m.Query(cmd, cb) } +func (m *QmpMonitor) MigrateCancel(cb StringCallback) { + m.HumanMonitorCommand("migrate_cancel", cb) +} + func (m *QmpMonitor) MigrateStartPostcopy(callback StringCallback) { var ( cmd = &Command{Execute: "migrate-start-postcopy"}