diff --git a/pkg/hostman/guestman/guesttasks.go b/pkg/hostman/guestman/guesttasks.go index caf5f17524..57f1d3758d 100644 --- a/pkg/hostman/guestman/guesttasks.go +++ b/pkg/hostman/guestman/guesttasks.go @@ -495,6 +495,7 @@ func (s *SGuestLiveMigrateTask) Start() { func (s *SGuestLiveMigrateTask) onSetZeroBlocks(res string) { if strings.Contains(strings.ToLower(res), "error") { + s.migrateTask = nil hostutils.TaskFailed(s.ctx, fmt.Sprintf("Migrate set capability zero-blocks error: %s", res)) return } @@ -502,6 +503,20 @@ func (s *SGuestLiveMigrateTask) onSetZeroBlocks(res string) { s.Monitor.MigrateSetCapability("auto-converge", "on", s.startMigrate) } +func (s *SGuestLiveMigrateTask) startRamMigrateTimeout() { + if !s.timeoutAt.IsZero() { + // timeout has been set + return + } + memMb, _ := s.Desc.Int("mem") + migSeconds := int(memMb) / options.HostOptions.MigrateExpectRate + if migSeconds < options.HostOptions.MinMigrateTimeoutSeconds { + migSeconds = options.HostOptions.MinMigrateTimeoutSeconds + } + s.timeoutAt = time.Now().Add(time.Second * time.Duration(migSeconds)) + log.Infof("migrate timeout seconds: %d now: %v expectfinial: %v", migSeconds, time.Now(), s.timeoutAt) +} + func (s *SGuestLiveMigrateTask) startMigrate(res string) { if strings.Contains(strings.ToLower(res), "error") { s.migrateTask = nil @@ -509,15 +524,9 @@ func (s *SGuestLiveMigrateTask) startMigrate(res string) { return } - memMb, _ := s.Desc.Int("mem") - migSeconds := int(memMb) / options.HostOptions.MigrateExpectRate - if migSeconds < options.HostOptions.MinMigrateTimeoutSeconds { - migSeconds = options.HostOptions.MinMigrateTimeoutSeconds - } - s.timeoutAt = time.Now().Add(time.Second * time.Duration(migSeconds)) - log.Infof("migrate timeout seconds: %d now: %v expectfinial: %v", migSeconds, time.Now(), s.timeoutAt) var copyIncremental = false if s.params.IsLocal { + // copy disk data copyIncremental = true } s.Monitor.Migrate(fmt.Sprintf("tcp:%s:%d", s.params.DestIp, s.params.DestPort), @@ -550,8 +559,12 @@ func (s *SGuestLiveMigrateTask) onGetMigrateStatus(status string) { s.migrateTask = nil close(s.c) hostutils.TaskFailed(s.ctx, fmt.Sprintf("Query migrate got status: %s", status)) - } else if !s.params.IsLocal && !s.doTimeoutMigrate { - if s.timeoutAt.Before(time.Now()) { + } else if status == "migrate_disk_copy" { + // do nothing, simply wait + } else if status == "migrate_ram_copy" { + if s.timeoutAt.IsZero() { + s.startRamMigrateTimeout() + } else if !s.doTimeoutMigrate && s.timeoutAt.Before(time.Now()) { log.Warningf("migrate timeout, force stop to finish migrate") // timeout, start memory postcopy // https://wiki.qemu.org/Features/PostCopyLiveMigration diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index e920070872..bd66a99c0c 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -629,13 +629,36 @@ func (m *QmpMonitor) GetMigrateStatus(callback StringCallback) { if res.ErrorVal != nil { callback(res.ErrorVal.Error()) } else { + /* + {"expected-downtime":300,"ram":{"dirty-pages-rate":0,"dirty-sync-count":1,"duplicate":2966538,"mbps":268.5672,"normal":148629,"normal-bytes":608784384,"page-size":4096,"postcopy-requests":0,"remaining":142815232,"skipped":0,"total":12902539264,"transferred":636674057},"setup-time":65,"status":"active","total-time":20002} + {"disk":{"dirty-pages-rate":0,"dirty-sync-count":0,"duplicate":0,"mbps":0,"normal":0,"normal-bytes":0,"page-size":0,"postcopy-requests":0,"remaining":0,"skipped":0,"total":139586437120,"transferred":139586437120},"expected-downtime":300,"ram":{"dirty-pages-rate":0,"dirty-sync-count":1,"duplicate":193281,"mbps":268.44264,"normal":62311,"normal-bytes":255225856,"page-size":4096,"postcopy-requests":0,"remaining":44474368,"skipped":0,"total":1091379200,"transferred":257555032},"setup-time":15,"status":"active","total-time":10002} + */ ret, err := jsonutils.Parse(res.Return) if err != nil { log.Errorf("Parse qmp res error %s: %s", m.server, err) callback("") } else { log.Infof("Query migrate status %s: %s", m.server, ret.String()) + status, _ := ret.GetString("status") + if status == "active" { + ramTotal, _ := ret.Int("ram", "total") + ramRemain, _ := ret.Int("ram", "remaining") + ramMbps, _ := ret.Float("ram", "mbps") + diskTotal, _ := ret.Int("disk", "total") + diskRemain, _ := ret.Int("disk", "remaining") + diskMbps, _ := ret.Float("disk", "mbps") + if diskRemain > 0 { + status = "migrate_disk_copy" + } else if ramRemain > 0 { + status = "migrate_ram_copy" + } + mbps := ramMbps + diskMbps + progress := (1 - float64(diskRemain+ramRemain)/float64(diskTotal+ramTotal)) * 100.0 + log.Debugf("progress: %f mbps: %f", progress, mbps) + hostutils.UpdateServerProgress(context.Background(), m.sid, progress, mbps) + } + callback(status) } }