From d81d6cc36832ccd0e3fde0c1c29119d4e8ac7efb Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Tue, 11 Jan 2022 02:48:17 +0800 Subject: [PATCH] fix: guest sync status persistent block_stream status --- pkg/compute/guestdrivers/baremetals.go | 4 +- pkg/compute/guestdrivers/container.go | 13 +-- pkg/compute/guestdrivers/esxi.go | 65 +++++++-------- pkg/compute/guestdrivers/kvm.go | 21 +++-- pkg/compute/guestdrivers/managedvirtual.go | 39 ++++----- pkg/compute/models/guest_actions.go | 9 ++- pkg/compute/models/guestdrivers.go | 2 +- pkg/compute/tasks/guest_syncstatus_task.go | 15 ++-- .../guestman/guesthandlers/guesthandler.go | 9 +-- pkg/hostman/guestman/guestman.go | 79 ++++++++++--------- pkg/hostman/guestman/libvirt.go | 2 +- pkg/hostman/guestman/qemu-kvm.go | 39 ++++++--- pkg/hostman/monitor/qmp.go | 1 + 13 files changed, 171 insertions(+), 127 deletions(-) diff --git a/pkg/compute/guestdrivers/baremetals.go b/pkg/compute/guestdrivers/baremetals.go index dead2c293f..ac5804344d 100644 --- a/pkg/compute/guestdrivers/baremetals.go +++ b/pkg/compute/guestdrivers/baremetals.go @@ -392,8 +392,8 @@ func (self *SBaremetalGuestDriver) RequestUndeployGuestOnHost(ctx context.Contex return err } -func (self *SBaremetalGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential) (jsonutils.JSONObject, error) { - return nil, fmt.Errorf("baremetal doesn't support RequestSyncstatusOnHost") +func (self *SBaremetalGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error { + return errors.Wrap(httperrors.ErrNotSupported, "baremetal doesn't support RequestSyncstatusOnHost") } func (self *SBaremetalGuestDriver) StartGuestSyncstatusTask(guest *models.SGuest, ctx context.Context, userCred mcclient.TokenCredential, parentTaskId string) error { diff --git a/pkg/compute/guestdrivers/container.go b/pkg/compute/guestdrivers/container.go index 164b5bbc23..cf58c10eac 100644 --- a/pkg/compute/guestdrivers/container.go +++ b/pkg/compute/guestdrivers/container.go @@ -116,11 +116,14 @@ func (self *SContainerDriver) RequestDetachDisk(ctx context.Context, guest *mode return containerUseKubectlError } -func (self *SContainerDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential) (jsonutils.JSONObject, error) { - // always return running - status := jsonutils.NewDict() - status.Add(jsonutils.NewString("running"), "status") - return status, nil +func (self *SContainerDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + // always return running + status := jsonutils.NewDict() + status.Add(jsonutils.NewString("running"), "status") + return status, nil + }) + return nil } func (self *SContainerDriver) CanKeepDetachDisk() bool { diff --git a/pkg/compute/guestdrivers/esxi.go b/pkg/compute/guestdrivers/esxi.go index 14405a73e6..f0c7897121 100644 --- a/pkg/compute/guestdrivers/esxi.go +++ b/pkg/compute/guestdrivers/esxi.go @@ -681,43 +681,46 @@ func (self *SESXiGuestDriver) RequestStopOnHost(ctx context.Context, guest *mode return nil } -func (self *SESXiGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential) (jsonutils.JSONObject, error) { - ihost, err := host.GetIHost() - if err != nil { - return nil, err - } - ivm, err := ihost.GetIVMById(guest.GetExternalId()) - if err != nil && errors.Cause(err) != errors.ErrNotFound { - return nil, err - } - // VM may be migrated by Vcenter, try to find VM from whole datacenter. - if err != nil { - ehost := ihost.(*esxi.SHost) - dc, err := ehost.GetDatacenter() +func (self *SESXiGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + ihost, err := host.GetIHost() if err != nil { - return nil, errors.Wrapf(err, "ehost.GetDatacenter") + return nil, errors.Wrap(err, "host.GetIHost") } - vm, err := dc.FetchVMById(guest.GetExternalId()) + ivm, err := ihost.GetIVMById(guest.GetExternalId()) if err != nil { - log.Errorf("fail to find ivm by id %q in dc %q: %v", guest.GetExternalId(), dc.GetName(), err) - return nil, err + if errors.Cause(err) != errors.ErrNotFound { + return nil, errors.Wrap(err, "ihost.GetIVMById") + } + // VM may be migrated by Vcenter, try to find VM from whole datacenter. + ehost := ihost.(*esxi.SHost) + dc, err := ehost.GetDatacenter() + if err != nil { + return nil, errors.Wrapf(err, "ehost.GetDatacenter") + } + vm, err := dc.FetchVMById(guest.GetExternalId()) + if err != nil { + log.Errorf("fail to find ivm by id %q in dc %q: %v", guest.GetExternalId(), dc.GetName(), err) + return nil, errors.Wrap(err, "dc.FetchVMById") + } + ihost = vm.GetIHost() + host = models.HostManager.FetchHostByExtId(ihost.GetGlobalId()) + if host == nil { + return nil, errors.Wrapf(errors.ErrNotFound, "find ivm %q in ihost %q which is not existed here", guest.GetExternalId(), ihost.GetGlobalId()) + } + ivm = vm } - ihost = vm.GetIHost() - host = models.HostManager.FetchHostByExtId(ihost.GetGlobalId()) - if host == nil { - return nil, errors.Wrapf(errors.ErrNotFound, "find ivm %q in ihost %q which is not existed here", guest.GetExternalId(), ihost.GetGlobalId()) + err = guest.SyncAllWithCloudVM(ctx, userCred, host, ivm, true) + if err != nil { + return nil, errors.Wrap(err, "guest.SyncAllWithCloudVM") } - ivm = vm - } - err = guest.SyncAllWithCloudVM(ctx, userCred, host, ivm, true) - if err != nil { - return nil, err - } - status := GetCloudVMStatus(ivm) - body := jsonutils.NewDict() - body.Add(jsonutils.NewString(status), "status") - return body, nil + status := GetCloudVMStatus(ivm) + body := jsonutils.NewDict() + body.Add(jsonutils.NewString(status), "status") + return body, nil + }) + return nil } func (self *SESXiGuestDriver) ValidateRebuildRoot(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, input *api.ServerRebuildRootInput) (*api.ServerRebuildRootInput, error) { diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index c2d646cb81..b3da7f2f1c 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -17,9 +17,9 @@ package guestdrivers import ( "context" "fmt" - "net/http" "strconv" "strings" + "time" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -304,17 +304,24 @@ func (self *SKVMGuestDriver) RequestStartOnHost(ctx context.Context, guest *mode return nil } -func (self *SKVMGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential) (jsonutils.JSONObject, error) { - header := http.Header{} - header.Set(mcclient.AUTH_TOKEN, userCred.GetTokenString()) - header.Set(mcclient.REGION_VERSION, "v2") +func (self *SKVMGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error { + header := self.getTaskRequestHeader(task) url := fmt.Sprintf("%s/servers/%s/status", host.ManagerUri, guest.Id) _, res, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "GET", url, header, nil, true) if err != nil { - return nil, err + return err } - return res, nil + statusStr, _ := res.GetString("status") + if len(statusStr) > 0 { + // may be an old version host, use sync request + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + // delay response to ensure event order + time.Sleep(time.Second) + return res, nil + }) + } + return nil } func (self *SKVMGuestDriver) OnDeleteGuestFinalCleanup(ctx context.Context, guest *models.SGuest, userCred mcclient.TokenCredential) error { diff --git a/pkg/compute/guestdrivers/managedvirtual.go b/pkg/compute/guestdrivers/managedvirtual.go index 4a298ed0f0..97783f7a56 100644 --- a/pkg/compute/guestdrivers/managedvirtual.go +++ b/pkg/compute/guestdrivers/managedvirtual.go @@ -818,26 +818,29 @@ func (self *SManagedVirtualizedGuestDriver) RequestStopOnHost(ctx context.Contex return nil } -func (self *SManagedVirtualizedGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential) (jsonutils.JSONObject, error) { - ihost, err := host.GetIHost() - if err != nil { - return nil, err - } - ivm, err := ihost.GetIVMById(guest.ExternalId) - if err != nil { - log.Errorf("fail to find ivm by id %s", err) - return nil, err - } +func (self *SManagedVirtualizedGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential, task taskman.ITask) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + ihost, err := host.GetIHost() + if err != nil { + return nil, errors.Wrap(err, "host.GetIHost") + } + ivm, err := ihost.GetIVMById(guest.ExternalId) + if err != nil { + log.Errorf("fail to find ivm by id %s", err) + return nil, errors.Wrap(err, "ihost.GetIVMById") + } - err = guest.SyncAllWithCloudVM(ctx, userCred, host, ivm, true) - if err != nil { - return nil, err - } + err = guest.SyncAllWithCloudVM(ctx, userCred, host, ivm, true) + if err != nil { + return nil, errors.Wrap(err, "guest.SyncAllWithCloudVM") + } - status := GetCloudVMStatus(ivm) - body := jsonutils.NewDict() - body.Add(jsonutils.NewString(status), "status") - return body, nil + status := GetCloudVMStatus(ivm) + body := jsonutils.NewDict() + body.Add(jsonutils.NewString(status), "status") + return body, nil + }) + return nil } func (self *SManagedVirtualizedGuestDriver) GetGuestVncInfo(ctx context.Context, userCred mcclient.TokenCredential, guest *models.SGuest, host *models.SHost, input *cloudprovider.ServerVncInput) (*cloudprovider.ServerVncOutput, error) { diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index c57d3dbb5c..5193a15895 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -2716,13 +2716,16 @@ func (self *SGuest) isNotRunningStatus(status string) bool { func (self *SGuest) PerformStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input apis.PerformStatusInput) (jsonutils.JSONObject, error) { preStatus := self.Status + + if len(self.BackupHostId) == 0 && input.Status == api.VM_RUNNING && input.BlockJobsCount > 0 { + input.Status = api.VM_BLOCK_STREAM + } _, err := self.SVirtualResourceBase.PerformStatus(ctx, userCred, query, input) if err != nil { - return nil, err + return nil, errors.Wrap(err, "SVirtualResourceBase.PerformStatus") } - status := input.Status - if len(self.BackupHostId) > 0 && status == api.VM_RUNNING && input.BlockJobsCount > 0 { + if len(self.BackupHostId) > 0 && input.Status == api.VM_RUNNING && input.BlockJobsCount > 0 { self.SetMetadata(ctx, api.MIRROR_JOB, api.MIRROR_JOB_READY, userCred) } else if ispId := self.GetMetadata(ctx, api.BASE_INSTANCE_SNAPSHOT_ID, userCred); len(ispId) > 0 { ispM, err := InstanceSnapshotManager.FetchById(ispId) diff --git a/pkg/compute/models/guestdrivers.go b/pkg/compute/models/guestdrivers.go index 990e0dfd91..c8b95ff349 100644 --- a/pkg/compute/models/guestdrivers.go +++ b/pkg/compute/models/guestdrivers.go @@ -105,7 +105,7 @@ type IGuestDriver interface { RequestSyncConfigOnHost(ctx context.Context, guest *SGuest, host *SHost, task taskman.ITask) error RequestSyncSecgroupsOnHost(ctx context.Context, guest *SGuest, host *SHost, task taskman.ITask) error - RequestSyncstatusOnHost(ctx context.Context, guest *SGuest, host *SHost, userCred mcclient.TokenCredential) (jsonutils.JSONObject, error) + RequestSyncstatusOnHost(ctx context.Context, guest *SGuest, host *SHost, userCred mcclient.TokenCredential, task taskman.ITask) error RequestStartOnHost(ctx context.Context, guest *SGuest, host *SHost, userCred mcclient.TokenCredential, task taskman.ITask) error diff --git a/pkg/compute/tasks/guest_syncstatus_task.go b/pkg/compute/tasks/guest_syncstatus_task.go index b1fd6fb051..22dd928f24 100644 --- a/pkg/compute/tasks/guest_syncstatus_task.go +++ b/pkg/compute/tasks/guest_syncstatus_task.go @@ -45,13 +45,13 @@ func (self *GuestSyncstatusTask) OnInit(ctx context.Context, obj db.IStandaloneM self.SetStageComplete(ctx, nil) return } - body, err := guest.GetDriver().RequestSyncstatusOnHost(ctx, guest, host, self.UserCred) + self.SetStage("OnGetStatusComplete", nil) + err := guest.GetDriver().RequestSyncstatusOnHost(ctx, guest, host, self.UserCred, self) if err != nil { log.Errorf("request_syncstatus_on_host: %s", err) - self.OnGetStatusFail(ctx, guest, err) + self.OnGetStatusCompleteFailed(ctx, guest, jsonutils.NewString(err.Error())) return } - self.OnGetStatusSucc(ctx, guest, body) } func (self *GuestSyncstatusTask) getOriginStatus() string { @@ -59,7 +59,9 @@ func (self *GuestSyncstatusTask) getOriginStatus() string { return os } -func (self *GuestSyncstatusTask) OnGetStatusSucc(ctx context.Context, guest *models.SGuest, body jsonutils.JSONObject) { +func (self *GuestSyncstatusTask) OnGetStatusComplete(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { + guest := obj.(*models.SGuest) + log.Debugf("OnGetStatusSucc guest %s(%s) status %s", guest.Name, guest.Id, body) statusStr, _ := body.GetString("status") switch statusStr { case cloudprovider.CloudVMStatusRunning: @@ -95,8 +97,9 @@ func (self *GuestSyncstatusTask) OnGetStatusSucc(ctx context.Context, guest *mod // logclient.AddActionLog(guest, logclient.ACT_VM_SYNC_STATUS, "", self.UserCred, true) } -func (self *GuestSyncstatusTask) OnGetStatusFail(ctx context.Context, guest *models.SGuest, err error) { - guest.SetStatus(self.UserCred, api.VM_UNKNOWN, err.Error()) +func (self *GuestSyncstatusTask) OnGetStatusCompleteFailed(ctx context.Context, obj db.IStandaloneModel, err jsonutils.JSONObject) { + guest := obj.(*models.SGuest) + guest.SetStatus(self.UserCred, api.VM_UNKNOWN, err.String()) self.SetStageComplete(ctx, nil) // logclient.AddActionLog(guest, logclient.ACT_VM_SYNC_STATUS, err, self.UserCred, false) } diff --git a/pkg/hostman/guestman/guesthandlers/guesthandler.go b/pkg/hostman/guestman/guesthandlers/guesthandler.go index d3b867f00a..b659261cbd 100644 --- a/pkg/hostman/guestman/guesthandlers/guesthandler.go +++ b/pkg/hostman/guestman/guesthandlers/guesthandler.go @@ -112,12 +112,9 @@ func guestActions(f actionFunc) appsrv.FilterHandler { func getStatus(ctx context.Context, w http.ResponseWriter, r *http.Request) { params, _, _ := appsrv.FetchEnv(ctx, w, r) - var status, blockJobsCount = guestman.GetGuestManager().StatusWithBlockJobsCount(params[""]) - res := map[string]interface{}{ - "status": status, - "block_jobs_count": blockJobsCount, - } - appsrv.SendStruct(w, res) + sid := params[""] + hostutils.DelayTaskWithoutReqctx(ctx, guestman.GetGuestManager().StatusWithBlockJobsCount, sid) + hostutils.ResponseOk(ctx, w) } func cpusetBalance(ctx context.Context, w http.ResponseWriter, r *http.Request) { diff --git a/pkg/hostman/guestman/guestman.go b/pkg/hostman/guestman/guestman.go index 2a67b699f0..84c056c524 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -624,47 +624,48 @@ func (m *SGuestManager) CpusetBalance(ctx context.Context, params interface{}) ( } func (m *SGuestManager) Status(sid string) string { - status := m.GetStatus(sid) - if status == GUEST_RUNNING { - guest, _ := m.GetServer(sid) - if guest.Monitor == nil && !guest.IsStopping() { - guest.StartMonitor(context.Background()) - } - } + status := m.getStatus(sid) return status } -func (m *SGuestManager) StatusWithBlockJobsCount(sid string) (string, int) { - status := m.GetStatus(sid) - blockJobsCount := 0 +func (m *SGuestManager) StatusWithBlockJobsCount(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { + sid := params.(string) + status := m.getStatus(sid) if status == GUEST_RUNNING { guest, _ := m.GetServer(sid) - if guest.Monitor == nil && !guest.IsStopping() { - guest.StartMonitor(context.Background()) + var runCb = func() { + if guest.IsMaster() { + mirrorStatus := guest.MirrorJobStatus() + if mirrorStatus.InProcess() { + status = GUEST_BLOCK_STREAM + } else if mirrorStatus.IsFailed() { + timeutils2.AddTimeout(1*time.Second, + func() { guest.SyncMirrorJobFailed("Block job missing") }) + status = GUEST_BLOCK_STREAM_FAIL + } + } + blockJobsCount := guest.BlockJobsCount() + body := jsonutils.NewDict() + body.Set("status", jsonutils.NewString(status)) + body.Set("block_jobs_count", jsonutils.NewInt(int64(blockJobsCount))) + hostutils.TaskComplete(ctx, body) } - blockJobsCount = guest.BlockJobsCount() + if guest.Monitor == nil && !guest.IsStopping() { + guest.StartMonitor(context.Background(), runCb) + } else { + runCb() + } + return nil, nil } - return status, blockJobsCount + body := jsonutils.NewDict() + body.Set("status", jsonutils.NewString(status)) + hostutils.TaskComplete(ctx, body) + return nil, nil } -func (m *SGuestManager) GetStatus(sid string) string { +func (m *SGuestManager) getStatus(sid string) string { if guest, ok := m.GetServer(sid); ok { - if guest.IsRunning() && guest.Monitor != nil && guest.IsMaster() { - mirrorStatus := guest.MirrorJobStatus() - if mirrorStatus.InProcess() { - return GUEST_BLOCK_STREAM - } else if mirrorStatus.IsFailed() { - timeutils2.AddTimeout(1*time.Second, - func() { guest.SyncMirrorJobFailed("Block job missing") }) - return GUEST_BLOCK_STREAM_FAIL - } else { - return GUEST_RUNNING - } - } if guest.IsRunning() { - if guest.Monitor != nil && guest.BlockJobsCount() > 0 { - return GUEST_BLOCK_STREAM - } return GUEST_RUNNING } else if guest.IsSuspend() { return GUEST_SUSPEND @@ -937,15 +938,19 @@ func (m *SGuestManager) Resume(ctx context.Context, sid string, isLiveMigrate bo if guest.IsStopping() || guest.IsStopped() { return nil, httperrors.NewInvalidStatusError("resume stopped server???") } + var cb = func() { + resumeTask := NewGuestResumeTask(ctx, guest, !isLiveMigrate) + if isLiveMigrate { + guest.StartPresendArp() + } + resumeTask.Start() + } if guest.Monitor == nil { - guest.StartMonitor(ctx) + guest.StartMonitor(ctx, cb) return nil, nil + } else { + cb() } - resumeTask := NewGuestResumeTask(ctx, guest, !isLiveMigrate) - if isLiveMigrate { - guest.StartPresendArp() - } - resumeTask.Start() return nil, nil } @@ -991,7 +996,7 @@ func (m *SGuestManager) CancelBlockJobs(ctx context.Context, params interface{}) if !ok { return nil, hostutils.ParamsError } - status := m.GetStatus(sid) + status := m.getStatus(sid) if status == GUSET_STOPPED { hostutils.TaskComplete(ctx, nil) return nil, nil diff --git a/pkg/hostman/guestman/libvirt.go b/pkg/hostman/guestman/libvirt.go index 8a55df061a..bc081b67e1 100644 --- a/pkg/hostman/guestman/libvirt.go +++ b/pkg/hostman/guestman/libvirt.go @@ -78,7 +78,7 @@ func (m *SGuestManager) GuestCreateFromLibvirt( if len(createConfig.MonitorPath) > 0 { if pid := findGuestProcessPid(guest.getOriginId(), "monitor.sock"); len(pid) > 0 { fileutils2.FilePutContents(guest.GetPidFilePath(), pid, false) - guest.StartMonitorWithImportGuestSocketFile(ctx, createConfig.MonitorPath) + guest.StartMonitorWithImportGuestSocketFile(ctx, createConfig.MonitorPath, nil) stopScript := guest.generateStopScript(nil) if err := fileutils2.FilePutContents(guest.GetStopScriptPath(), stopScript, false); err != nil { return nil, fmt.Errorf("Save stop script error %s", err) diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index b0a515e140..4f8378a762 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -283,7 +283,7 @@ func (s *SKVMGuestInstance) asyncScriptStart(ctx context.Context, params interfa if isStarted { log.Infof("Async start server %s success!", s.GetName()) s.syncMeta = s.CleanImportMetadata() - s.StartMonitor(ctx) + s.StartMonitor(ctx, nil) return nil, nil } log.Infof("Async start server %s failed: %s!!!", s.GetName(), err) @@ -339,7 +339,7 @@ func (s *SKVMGuestInstance) ImportServer(pendingDelete bool) { if s.IsRunning() { log.Infof("%s is running, pending_delete=%t", s.GetName(), pendingDelete) if !pendingDelete { - s.StartMonitor(context.Background()) + s.StartMonitor(context.Background(), nil) } } else { var action = "stopped" @@ -429,21 +429,28 @@ func (s *SKVMGuestInstance) GetMonitorPath() string { return monitorPath } -func (s *SKVMGuestInstance) StartMonitorWithImportGuestSocketFile(ctx context.Context, socketFile string) { +func (s *SKVMGuestInstance) StartMonitorWithImportGuestSocketFile(ctx context.Context, socketFile string, cb func()) { timeutils2.AddTimeout(100*time.Millisecond, func() { - s.Monitor = monitor.NewQmpMonitor( + var mon monitor.Monitor + mon = monitor.NewQmpMonitor( s.GetName(), s.Id, s.onImportGuestMonitorDisConnect, // on monitor disconnect func(err error) { s.onImportGuestMonitorTimeout(ctx, err) }, // on monitor timeout - func() { s.onImportGuestMonitorConnected(ctx) }, // on monitor connected + func() { + s.Monitor = mon + s.onImportGuestMonitorConnected(ctx) + if cb != nil { + cb() + } + }, // on monitor connected s.onReceiveQMPEvent, // on reveive qmp event ) - s.Monitor.ConnectWithSocket(socketFile) + mon.ConnectWithSocket(socketFile) }) } -func (s *SKVMGuestInstance) StartMonitor(ctx context.Context) { +func (s *SKVMGuestInstance) StartMonitor(ctx context.Context, cb func()) { if options.HostOptions.EnableQmpMonitor && s.GetQmpMonitorPort(-1) > 0 { // try qmp first, if qmp connect failed, use hmp timeutils2.AddTimeout(100*time.Millisecond, func() { @@ -456,6 +463,9 @@ func (s *SKVMGuestInstance) StartMonitor(ctx context.Context) { func() { s.Monitor = mon s.onMonitorConnected(ctx) + if cb != nil { + cb() + } }, // on monitor connected s.onReceiveQMPEvent, // on reveive qmp event ) @@ -467,7 +477,13 @@ func (s *SKVMGuestInstance) StartMonitor(ctx context.Context) { s.Id, s.onMonitorDisConnect, // on monitor disconnect func(err error) { s.onMonitorTimeout(ctx, err) }, // on monitor timeout - func() { s.onMonitorConnected(ctx) }, // on monitor connected + func() { + s.Monitor = mon + s.onMonitorConnected(ctx) + if cb != nil { + cb() + } + }, // on monitor connected ) err = mon.Connect("127.0.0.1", s.GetHmpMonitorPort(-1)) if err != nil { @@ -477,7 +493,7 @@ func (s *SKVMGuestInstance) StartMonitor(ctx context.Context) { } }) } else if monitorPath := s.GetMonitorPath(); len(monitorPath) > 0 { - s.StartMonitorWithImportGuestSocketFile(ctx, monitorPath) + s.StartMonitorWithImportGuestSocketFile(ctx, monitorPath, cb) } else { log.Errorf("Guest start monitor failed, can't get qmp monitor port or monitor path") } @@ -693,13 +709,16 @@ func (s *SKVMGuestInstance) MirrorJobStatus() MirrorJob { func (s *SKVMGuestInstance) BlockJobsCount() int { res := make(chan []monitor.BlockJob) + log.Debugf("BlockJobsCount start...") s.Monitor.GetBlockJobs(func(jobs []monitor.BlockJob) { res <- jobs }) select { - case <-time.After(time.Second * 3): + case <-time.After(time.Second * 30): + log.Debugf("BlockJobsCount timeout") return -1 case v := <-res: + log.Debugf("BlockJobsCount %d", len(v)) return len(v) } } diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index a2d7bfd84e..87cf69e823 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -702,6 +702,7 @@ func (m *QmpMonitor) blockJobs(res *Response) ([]BlockJob, error) { if err != nil { return nil, errors.Wrapf(err, "GetBlockJobs for %s parse %s", m.server, res.Return) } + log.Debugf("blockJobs response %s", ret) jobs := []BlockJob{} ret.Unmarshal(&jobs) defer func() {