diff --git a/pkg/compute/guestdrivers/base.go b/pkg/compute/guestdrivers/base.go index 93aba361e2..8d44e0a9d0 100644 --- a/pkg/compute/guestdrivers/base.go +++ b/pkg/compute/guestdrivers/base.go @@ -281,6 +281,10 @@ func (self *SBaseGuestDriver) RequestSyncToBackup(ctx context.Context, guest *mo return fmt.Errorf("Not Implement") } +func (self *SBaseGuestDriver) RequestSlaveBlockStreamDisks(ctx context.Context, guest *models.SGuest, task taskman.ITask) error { + return fmt.Errorf("Not Implement") +} + func (self *SBaseGuestDriver) GetMaxSecurityGroupCount() int { return 5 } diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index 73c0457f8c..58d9bab38c 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -695,6 +695,15 @@ func (self *SKVMGuestDriver) RequestSyncToBackup(ctx context.Context, guest *mod return nil } +func (self *SKVMGuestDriver) RequestSlaveBlockStreamDisks(ctx context.Context, guest *models.SGuest, task taskman.ITask) error { + host := models.HostManager.FetchHostById(guest.BackupHostId) + body := jsonutils.NewDict() + url := fmt.Sprintf("%s/servers/%s/slave-block-stream-disks", host.ManagerUri, guest.Id) + header := self.getTaskRequestHeader(task) + _, _, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, body, false) + return err +} + // kvm guest must add cpu first // if body has add_cpu_failed indicate dosen't exec add mem // 1. cpu added part of request --> add_cpu_failed: true && added_cpu: count diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 453c564855..86cf8607f1 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -2869,24 +2869,13 @@ func (self *SGuest) PerformStatus(ctx context.Context, userCred mcclient.TokenCr } 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, errors.Wrap(err, "SVirtualResourceBase.PerformStatus") } if self.HasBackupGuest() { - if input.Status == api.VM_RUNNING { - if err := self.TrySetGuestBackupMirrorJobReady(ctx, userCred); err != nil { - return nil, errors.Wrap(err, "set guest backup mirror job status ready") - } - } else if input.Status == api.VM_BLOCK_STREAM { - if err := self.SetGuestBackupMirrorJobInProgress(ctx, userCred); err != nil { - return nil, errors.Wrap(err, "set guest backup mirror job status inprogress") - } - } else if input.Status == api.VM_READY { + if input.Status == api.VM_READY { if err := self.ResetGuestQuorumChildIndex(ctx, userCred); err != nil { return nil, errors.Wrap(err, "reset guest quorum child index") } @@ -3421,6 +3410,16 @@ func (self *SGuest) PerformBlockStreamFailed(ctx context.Context, userCred mccli return nil, nil } +func (self *SGuest) PerformSlaveBlockStreamReady(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + if len(self.BackupHostId) > 0 { + if err := self.TrySetGuestBackupMirrorJobReady(ctx, userCred); err != nil { + return nil, errors.Wrap(err, "set guest backup mirror job status ready") + } + self.SetBackupGuestStatus(userCred, api.VM_RUNNING, "perform slave block stream ready") + } + return nil, nil +} + func (self *SGuest) PerformBlockMirrorReady(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { if self.Status == api.VM_BLOCK_STREAM || self.Status == api.VM_RUNNING { diskId, err := data.GetString("disk_id") diff --git a/pkg/compute/models/guestdrivers.go b/pkg/compute/models/guestdrivers.go index 99d31f97e7..d83a0f50ed 100644 --- a/pkg/compute/models/guestdrivers.go +++ b/pkg/compute/models/guestdrivers.go @@ -177,6 +177,7 @@ type IGuestDriver interface { RequestDeleteSnapshot(ctx context.Context, guest *SGuest, task taskman.ITask, params *jsonutils.JSONDict) error RequestReloadDiskSnapshot(ctx context.Context, guest *SGuest, task taskman.ITask, params *jsonutils.JSONDict) error RequestSyncToBackup(ctx context.Context, guest *SGuest, task taskman.ITask) error + RequestSlaveBlockStreamDisks(ctx context.Context, guest *SGuest, task taskman.ITask) error IsSupportEip() bool IsSupportPublicIp() bool diff --git a/pkg/compute/tasks/guest_backup_tasks.go b/pkg/compute/tasks/guest_backup_tasks.go index a8ae945773..9f35ddce98 100644 --- a/pkg/compute/tasks/guest_backup_tasks.go +++ b/pkg/compute/tasks/guest_backup_tasks.go @@ -86,10 +86,7 @@ func (self *GuestSwitchToBackupTask) OnBackupGuestStoped(ctx context.Context, gu self.OnFail(ctx, guest, jsonutils.NewString(fmt.Sprintf("Switch to backup guest error: %s", err))) return } - if err := guest.SetGuestBackupMirrorJobNotReady(ctx, self.UserCred); err != nil { - self.OnFail(ctx, guest, jsonutils.NewString("guest set metadata failed")) - return - } + db.OpsLog.LogEvent(guest, db.ACT_SWITCHED, "Switch to backup", self.UserCred) logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_SWITCH_TO_BACKUP, "Switch to backup", self.UserCred, true) oldStatus, _ := self.Params.GetString("old_status") @@ -197,7 +194,16 @@ func (self *GuestStartAndSyncToBackupTask) OnStartBackupGuestFailed(ctx context. func (self *GuestStartAndSyncToBackupTask) OnRequestSyncToBackup(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { guest.SetGuestBackupMirrorJobInProgress(ctx, self.UserCred) - guest.SetStatus(self.UserCred, api.VM_BLOCK_STREAM, "OnSyncToBackup") + err := guest.GetDriver().RequestSlaveBlockStreamDisks(ctx, guest, self) + if err != nil { + guest.SetGuestBackupMirrorJobFailed(ctx, self.UserCred) + guest.SetBackupGuestStatus(self.UserCred, api.VM_BLOCK_STREAM_FAIL, err.Error()) + self.SetStageFailed(ctx, jsonutils.NewString(err.Error())) + return + } + + guest.SetGuestBackupMirrorJobInProgress(ctx, self.UserCred) + guest.SetBackupGuestStatus(self.UserCred, api.VM_BLOCK_STREAM, "OnSyncToBackup") self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/ha_guest_start_task.go b/pkg/compute/tasks/ha_guest_start_task.go index a5f8bc25db..db4766a3dd 100644 --- a/pkg/compute/tasks/ha_guest_start_task.go +++ b/pkg/compute/tasks/ha_guest_start_task.go @@ -40,6 +40,7 @@ func (self *HAGuestStartTask) OnInit( guest := obj.(*models.SGuest) host := models.HostManager.FetchHostById(guest.BackupHostId) if host.HostStatus != api.HOST_ONLINE { + guest.SetGuestBackupMirrorJobFailed(ctx, self.UserCred) // request start master guest self.GuestStartTask.OnInit(ctx, guest, nil) } else { @@ -77,6 +78,15 @@ func (self *HAGuestStartTask) RequestStartBacking(ctx context.Context, guest *mo host := models.HostManager.FetchHostById(guest.BackupHostId) guest.SetStatus(self.UserCred, api.VM_BACKUP_STARTING, "") + if !guest.IsGuestBackupMirrorJobReady(ctx, self.UserCred) { + hostMaster := models.HostManager.FetchHostById(guest.HostId) + self.Params.Set("block_ready", jsonutils.JSONFalse) + diskUri := fmt.Sprintf("%s/disks", hostMaster.GetFetchUrl(true)) + self.Params.Set("disk_uri", jsonutils.NewString(diskUri)) + } else { + self.Params.Set("block_ready", jsonutils.JSONTrue) + } + err := guest.GetDriver().RequestStartOnHost(ctx, guest, host, self.UserCred, self) if err != nil { self.OnStartCompleteFailed(ctx, guest, jsonutils.NewString(err.Error())) @@ -113,3 +123,17 @@ func (self *HAGuestStartTask) OnStartBackupGuestCompleteFailed( guest.SetBackupGuestStatus(self.UserCred, api.VM_START_FAILED, data.String()) self.OnStartCompleteFailed(ctx, guest, data) } + +func (self *HAGuestStartTask) OnStartComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + if !guest.IsGuestBackupMirrorJobReady(ctx, self.UserCred) { + if err := guest.GetDriver().RequestSlaveBlockStreamDisks(ctx, guest, self); err != nil { + guest.SetGuestBackupMirrorJobFailed(ctx, self.UserCred) + guest.SetBackupGuestStatus(self.UserCred, api.VM_BLOCK_STREAM_FAIL, err.Error()) + } else { + guest.SetGuestBackupMirrorJobInProgress(ctx, self.UserCred) + guest.SetBackupGuestStatus(self.UserCred, api.VM_BLOCK_STREAM, "on RequestSlaveBlockStreamDisks") + } + } + + self.GuestStartTask.OnStartComplete(ctx, guest, data) +} diff --git a/pkg/hostman/guestman/guesthandlers/guesthandler.go b/pkg/hostman/guestman/guesthandlers/guesthandler.go index 60daba2e91..e63df5465d 100644 --- a/pkg/hostman/guestman/guesthandlers/guesthandler.go +++ b/pkg/hostman/guestman/guesthandlers/guesthandler.go @@ -79,6 +79,7 @@ func AddGuestTaskHandler(prefix string, app *appsrv.Application) { "live-migrate": guestLiveMigrate, "resume": guestResume, "block-replication": guestBlockReplication, + "slave-block-stream-disks": slaveGuestBlockStreamDisks, "hotplug-cpu-mem": guestHotplugCpuMem, "cancel-block-jobs": guestCancelBlockJobs, "cancel-block-replication": guestCancelBlockReplication, @@ -501,6 +502,14 @@ func guestBlockReplication(ctx context.Context, userCred mcclient.TokenCredentia return nil, nil } +func slaveGuestBlockStreamDisks(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) + } + return nil, guest.SlaveDisksBlockStream() +} + func guestCancelBlockJobs(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/guestman.go b/pkg/hostman/guestman/guestman.go index 738f896ff0..5399278633 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -750,20 +750,11 @@ func (m *SGuestManager) StatusWithBlockJobsCount(ctx context.Context, params int } else if status == GUEST_RUNNING { var runCb = func() { body := jsonutils.NewDict() - 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("drive-mirror job failed") }) - status = GUEST_BLOCK_STREAM_FAIL - } - body.Set("block_jobs_count", jsonutils.NewInt(int64(mirrorStatus.blockJobsCount))) - } else { - blockJobsCount := guest.BlockJobsCount() - body.Set("block_jobs_count", jsonutils.NewInt(int64(blockJobsCount))) + blockJobsCount := guest.BlockJobsCount() + if blockJobsCount > 0 { + status = GUEST_BLOCK_STREAM } + body.Set("block_jobs_count", jsonutils.NewInt(int64(blockJobsCount))) body.Set("status", jsonutils.NewString(status)) hostutils.TaskComplete(ctx, body) } @@ -1254,7 +1245,7 @@ func (m *SGuestManager) StartBlockReplication(ctx context.Context, params interf } hostutils.TaskComplete(ctx, nil) } - task := NewGuestBlockReplicationTask(ctx, guest, nbdOpts[1], nbdOpts[2], "top", onSucc, nil) + task := NewGuestBlockReplicationTask(ctx, guest, nbdOpts[1], nbdOpts[2], "full", onSucc, nil) task.Start() return nil, nil } diff --git a/pkg/hostman/guestman/guesttasks.go b/pkg/hostman/guestman/guesttasks.go index a1b7ad6d22..fddd199508 100644 --- a/pkg/hostman/guestman/guesttasks.go +++ b/pkg/hostman/guestman/guesttasks.go @@ -1961,6 +1961,52 @@ func (s *SDriveMirrorTask) startMirror(res string) { } } +type SDriveBackupTask struct { + *SKVMGuestInstance + + ctx context.Context + nbdUri string + onSucc func() + syncMode string + index int +} + +func NewDriveBackupTask( + ctx context.Context, s *SKVMGuestInstance, nbdUri, syncMode string, onSucc func(), +) *SDriveBackupTask { + return &SDriveBackupTask{ + SKVMGuestInstance: s, + ctx: ctx, + nbdUri: nbdUri, + syncMode: syncMode, + onSucc: onSucc, + } +} + +func (s *SDriveBackupTask) Start() { + s.startBackup("") +} + +func (s *SDriveBackupTask) startBackup(res string) { + log.Infof("drive backup results:%s", res) + if len(res) > 0 { + hostutils.TaskFailed(s.ctx, res) + return + } + disks := s.Desc.Disks + if s.index < len(disks) { + target := fmt.Sprintf("%s:exportname=drive_%d_backend", s.nbdUri, s.index) + s.Monitor.DriveBackup(s.startBackup, fmt.Sprintf("drive_%d", s.index), target, s.syncMode, "raw") + s.index += 1 + } else { + if s.onSucc != nil { + s.onSucc() + } else { + hostutils.TaskComplete(s.ctx, nil) + } + } +} + type SGuestBlockReplicationTask struct { *SKVMGuestInstance @@ -2015,7 +2061,11 @@ func (s *SGuestBlockReplicationTask) onXBlockdevChange(res string) { }, s.onNbdDriveAddSucc(drive, node)) s.index += 1 } else { - s.startDriveMirror() + if s.onSucc != nil { + s.onSucc() + } else { + hostutils.TaskComplete(s.ctx, nil) + } } } @@ -2035,11 +2085,6 @@ func (s *SGuestBlockReplicationTask) onNbdDriveAddSucc(parent, node string) moni } } -func (s *SGuestBlockReplicationTask) startDriveMirror() { - NewDriveMirrorTask(s.ctx, s.SKVMGuestInstance, - fmt.Sprintf("nbd:%s:%s", s.nbdHost, s.nbdPort), s.syncMode, true, s.onSucc).Start() -} - /** * GuestOnlineResizeDiskTask **/ @@ -2397,7 +2442,9 @@ func (task *CancelBlockReplication) Start() { }) }) } - task.SCancelBlockJobs.Start() + if task.ctx != nil { + hostutils.TaskComplete(task.ctx, nil) + } } type SCancelBlockJobs struct { diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index 09339d67f1..7bd21d6c69 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -518,7 +518,6 @@ func (s *SKVMGuestInstance) asyncScriptStart(ctx context.Context, params interfa migratePortInt64 := int64(migratePort) s.LiveMigrateDestPort = &migratePortInt64 } - data.Set("script_start", jsonutils.JSONTrue) err = s.saveScripts(data) if err != nil { @@ -840,51 +839,61 @@ func (s *SKVMGuestInstance) eventBlockJobReady(event *monitor.Event) { log.Errorf("block job missing event type") return } - // only dealwith event type mirror stype, _ := itype.(string) - if stype != "mirror" { + if stype != "mirror" && stype != "stream" { + return + } + iDevice, ok := event.Data["device"] + if !ok { + return + } + device := iDevice.(string) + if !strings.HasPrefix(device, "drive_") { + return + } + disks := s.Desc.Disks + diskIndex, err := strconv.Atoi(device[len("drive_"):]) + if err != nil || diskIndex < 0 || diskIndex >= len(disks) { + log.Errorf("failed get disk from index %d", diskIndex) + return + } + var diskId, diskPath string + for i := 0; i < len(disks); i++ { + index := disks[i].Index + if index == int8(diskIndex) { + diskId = disks[i].DiskId + diskPath = disks[i].Path + } + } + if len(diskId) == 0 { + log.Errorf("failed find disk %s", device) return } - if s.IsMaster() { // has backup server - mirrorStatus := s.MirrorJobStatus() - if mirrorStatus.IsSucc() { + if s.IsSlave() { // is backup server + disk, err := storageman.GetManager().GetDiskByPath(diskPath) + if err != nil { + log.Errorf("eventBlockJobReady failed get disk %s", diskPath) + return + } + disk.PostCreateFromImageFuse() + blockJobCount := s.BlockJobsCount() + if blockJobCount == 0 { for { - statusInput := &apis.PerformStatusInput{ - Status: api.VM_RUNNING, - Reason: "block job ready", - BlockJobsCount: mirrorStatus.BlockJobsCount(), - PowerStates: s.GetPowerStates(), - } - _, err := hostutils.UpdateServerStatus(context.Background(), s.GetId(), statusInput) + _, err := modules.Servers.PerformAction( + hostutils.GetComputeSession(context.Background()), s.GetId(), "slave-block-stream-ready", nil, + ) if err != nil { - log.Errorf("onReceiveQMPEvent update server status error: %s", err) + log.Errorf("onReceiveQMPEvent sync slave block stream ready error: %s", err) time.Sleep(3 * time.Second) } else { break } } - } else if mirrorStatus.IsFailed() { - s.SyncMirrorJobFailed("drive-mirror job failed") } } else { - iDevice, ok := event.Data["device"] - if !ok { - return - } - device := iDevice.(string) - if !strings.HasPrefix(device, "drive_") { - return - } - disks := s.Desc.Disks - log.Infof("mirror job ready disk index %s", device[len("drive_"):]) - diskIndex, err := strconv.Atoi(device[len("drive_"):]) - if err != nil || diskIndex < 0 || diskIndex >= len(disks) { - log.Errorf("failed get disk from index %d", diskIndex) - return - } params := jsonutils.NewDict() - params.Set("disk_id", jsonutils.NewString(disks[diskIndex].DiskId)) + params.Set("disk_id", jsonutils.NewString(diskId)) _, err = modules.Servers.PerformAction( hostutils.GetComputeSession(context.Background()), s.GetId(), "block-mirror-ready", params, @@ -1168,7 +1177,7 @@ func (s *SKVMGuestInstance) startDiskBackupMirror(ctx context.Context) { s.SyncMirrorJobFailed(res) s.DoResumeTask(ctx, true) } - NewGuestBlockReplicationTask(ctx, s, nbdOpts[1], nbdOpts[2], "top", onSucc, onFail).Start() + NewGuestBlockReplicationTask(ctx, s, nbdOpts[1], nbdOpts[2], "full", onSucc, onFail).Start() } } @@ -1187,9 +1196,27 @@ func (s *SKVMGuestInstance) startQemuBuiltInNbdServer(ctx context.Context) { } } s.Monitor.StartNbdServer(nbdServerPort, true, true, onNbdServerStarted) + } else { + s.SyncStatus("") } } +func (s *SKVMGuestInstance) SlaveDisksBlockStream() error { + errChan := make(chan string, 1) + disks := s.Desc.Disks + for i := 0; i < len(disks); i++ { + diskIndex := disks[i].Index + drive := fmt.Sprintf("drive_%d", diskIndex) + s.Monitor.BlockStream(drive, 0, 0, func(res string) { + errChan <- res + }) + if errStr := <-errChan; len(errStr) > 0 { + return fmt.Errorf("block stream disk %s: %s", drive, errStr) + } + } + return nil +} + func (s *SKVMGuestInstance) clearCgroup(pid int) { if pid == 0 && s.cgroupPid > 0 { pid = s.cgroupPid @@ -1217,58 +1244,6 @@ func (s *SKVMGuestInstance) DiskCount() int { return len(s.Desc.Disks) } -type MirrorJob struct { - mirrorJobStatus int - blockJobsCount int -} - -func (ms MirrorJob) IsSucc() bool { - return ms.mirrorJobStatus == 1 -} - -func (ms MirrorJob) IsFailed() bool { - return ms.mirrorJobStatus == -1 -} - -func (ms MirrorJob) InProcess() bool { - return ms.mirrorJobStatus == 0 -} - -func (ms MirrorJob) BlockJobsCount() int { - return ms.blockJobsCount -} - -func (s *SKVMGuestInstance) MirrorJobStatus() MirrorJob { - res := make(chan []monitor.BlockJob) - s.Monitor.GetBlockJobs(func(jobs []monitor.BlockJob) { - res <- jobs - }) - select { - case <-time.After(time.Second * 3): - return MirrorJob{0, -1} - case v := <-res: - mirrorJobCount := 0 - failedJobCount := 0 - for _, job := range v { - if job.Type != "mirror" { - continue - } - if job.IoStatus != "ok" { - failedJobCount += 1 - } - mirrorJobCount += 1 - } - if failedJobCount > 0 { - return MirrorJob{-1, len(v)} - } - if mirrorJobCount == 0 { - return MirrorJob{1, len(v)} - } else { - return MirrorJob{0, len(v)} - } - } -} - func (s *SKVMGuestInstance) BlockJobsCount() int { res := make(chan []monitor.BlockJob) log.Debugf("BlockJobsCount start...") @@ -1388,23 +1363,14 @@ func (s *SKVMGuestInstance) CheckBlockOrRunning(jobs int) { var status = api.VM_RUNNING if jobs > 0 { - if s.IsMaster() { - mirrorStatus := s.MirrorJobStatus() - if mirrorStatus.InProcess() { - status = api.VM_BLOCK_STREAM - } else if mirrorStatus.IsFailed() { - status = api.VM_BLOCK_STREAM_FAIL - s.SyncMirrorJobFailed("drive-mirror job failed") - } - } else { - // TODO: check block jobs ready - status = api.VM_BLOCK_STREAM - } + // TODO: check block jobs ready + status = api.VM_BLOCK_STREAM } var statusInput = &apis.PerformStatusInput{ Status: status, BlockJobsCount: jobs, PowerStates: s.GetPowerStates(), + IsSlave: s.IsSlave(), } _, err := hostutils.UpdateServerStatus(context.Background(), s.Id, statusInput) if err != nil { diff --git a/pkg/hostman/guestman/qemu-kvmhelper.go b/pkg/hostman/guestman/qemu-kvmhelper.go index 5223c088a8..995b3aa3db 100644 --- a/pkg/hostman/guestman/qemu-kvmhelper.go +++ b/pkg/hostman/guestman/qemu-kvmhelper.go @@ -15,6 +15,7 @@ package guestman import ( + "context" "fmt" "net" "path" @@ -39,9 +40,7 @@ import ( "yunion.io/x/onecloud/pkg/hostman/monitor" "yunion.io/x/onecloud/pkg/hostman/options" "yunion.io/x/onecloud/pkg/hostman/storageman" - "yunion.io/x/onecloud/pkg/util/fileutils2" "yunion.io/x/onecloud/pkg/util/procutils" - "yunion.io/x/onecloud/pkg/util/qemuimg" "yunion.io/x/onecloud/pkg/util/qemutils" ) @@ -509,8 +508,12 @@ function nic_mtu() { input.LiveMigratePort = uint(*s.LiveMigrateDestPort) } - if s.Desc.IsSlave && jsonutils.QueryBoolean(data, "script_start", false) { - if err := s.slaveDiskPrepare(input); err != nil { + if s.Desc.IsSlave && !jsonutils.QueryBoolean(data, "block_ready", false) { + diskUri, err := data.GetString("disk_uri") + if err != nil { + return "", errors.Wrap(err, "guest start missing disk uri") + } + if err := s.slaveDiskPrepare(input, diskUri); err != nil { return "", err } } @@ -533,44 +536,19 @@ echo $CMD` return cmd, nil } -func (s *SKVMGuestInstance) slaveDiskPrepare(input *qemu.GenerateStartOptionsInput) error { +func (s *SKVMGuestInstance) slaveDiskPrepare(input *qemu.GenerateStartOptionsInput, diskUri string) error { for i := 0; i < len(input.GuestDesc.Disks); i++ { diskPath := input.GuestDesc.Disks[i].Path d, err := storageman.GetManager().GetDiskByPath(diskPath) if err != nil { return errors.Wrapf(err, "GetDiskByPath(%s)", diskPath) } - disk, err := qemuimg.NewQemuImage(d.GetPath()) - if err != nil { - return errors.Wrapf(err, "qemuimg.NewQemuImage(%s)", diskPath) + if output, err := procutils.NewCommand("rm", "-f", diskPath).Output(); err != nil { + return errors.Errorf("failed delete slave top disk file %s %s", output, err) } - backendPath := diskPath + ".backend" - if !fileutils2.Exists(backendPath) { - output, err := procutils.NewCommand("mv", "-f", diskPath, backendPath).Output() - if err != nil { - return errors.Wrapf(err, "mv %s to %s failed %s", diskPath, backendPath, output) - } - diskTop, err := qemuimg.NewQemuImage(diskPath) - if err != nil { - return errors.Wrap(err, "qemuimg.NewQemuImage") - } - if err = diskTop.CreateQcow2(0, false, backendPath, "", "", ""); err != nil { - return errors.Wrap(err, "create qcow2") - } - } else { - if disk.BackFilePath != backendPath { - return errors.Errorf("backend file %s exist but not a backing file", backendPath) - } - if output, err := procutils.NewCommand("rm", "-f", diskPath).Output(); err != nil { - return errors.Errorf("failed delete slave top disk file %s %s", output, err) - } - diskTop, err := qemuimg.NewQemuImage(diskPath) - if err != nil { - return errors.Wrap(err, "qemuimg.NewQemuImage") - } - if err = diskTop.CreateQcow2(0, false, backendPath, "", "", ""); err != nil { - return errors.Wrap(err, "create qcow2") - } + diskUrl := fmt.Sprintf("%s/%s", diskUri, input.GuestDesc.Disks[i].DiskId) + if err := d.CreateFromImageFuse(context.Background(), diskUrl, 0, nil); err != nil { + return errors.Wrap(err, "failed create slave disk") } } return nil diff --git a/pkg/hostman/guestman/qemu/generate.go b/pkg/hostman/guestman/qemu/generate.go index 3820d5f980..b3ff7e1cd3 100644 --- a/pkg/hostman/guestman/qemu/generate.go +++ b/pkg/hostman/guestman/qemu/generate.go @@ -218,10 +218,7 @@ func generateDisksOptions(drvOpt QemuOptions, disks []*desc.SGuestDisk, isEncryp if isMaster { opts = append(opts, getMasterDiskDriveOption(drvOpt, disk, isEncrypt)) } else { - opts = append(opts, getDiskDriveOption(drvOpt, disk, isEncrypt, false)) - if isSlave { // append slave backend disk - opts = append(opts, getDiskDriveOption(drvOpt, disk, isEncrypt, true)) - } + opts = append(opts, getDiskDriveOption(drvOpt, disk, isEncrypt)) } opts = append(opts, getDiskDeviceOption(drvOpt, disk)) } @@ -246,22 +243,16 @@ func getMasterDiskDriveOption(drvOpt QemuOptions, disk *desc.SGuestDisk, isEncry return drvOpt.Drive(opt) } -func getDiskDriveOption(drvOpt QemuOptions, disk *desc.SGuestDisk, isEncrypt, isSlave bool) string { +func getDiskDriveOption(drvOpt QemuOptions, disk *desc.SGuestDisk, isEncrypt bool) string { format := disk.Format diskIndex := disk.Index cacheMode := disk.CacheMode aioMode := disk.AioMode var opt string - if isSlave { - opt = fmt.Sprintf("file=$DISK_%d.backend", diskIndex) - opt += ",if=none" - opt += fmt.Sprintf(",id=drive_%d_backend", diskIndex) - } else { - opt = fmt.Sprintf("file=$DISK_%d", diskIndex) - opt += ",if=none" - opt += fmt.Sprintf(",id=drive_%d", diskIndex) - } + opt = fmt.Sprintf("file=$DISK_%d", diskIndex) + opt += ",if=none" + opt += fmt.Sprintf(",id=drive_%d", diskIndex) if len(format) == 0 || format == "qcow2" { // pass # qemu will automatically detect image format diff --git a/pkg/hostman/monitor/hmp.go b/pkg/hostman/monitor/hmp.go index 74984e636d..bb24f42014 100644 --- a/pkg/hostman/monitor/hmp.go +++ b/pkg/hostman/monitor/hmp.go @@ -414,6 +414,15 @@ func (m *HmpMonitor) DriveMirror(callback StringCallback, drive, target, syncMod m.Query(cmd, callback) } +func (m *HmpMonitor) DriveBackup(callback StringCallback, drive, target, syncMode, format string) { + cmd := "drive_backup -n" + if syncMode == "full" { + cmd += " -f" + } + cmd += fmt.Sprintf(" %s %s %s", drive, target, format) + m.Query(cmd, callback) +} + func (m *HmpMonitor) BlockStream(drive string, _, _ int, callback StringCallback) { var ( speed = 500 // limit 500 MB/s diff --git a/pkg/hostman/monitor/monitor.go b/pkg/hostman/monitor/monitor.go index c8f16ec010..d192232347 100644 --- a/pkg/hostman/monitor/monitor.go +++ b/pkg/hostman/monitor/monitor.go @@ -224,6 +224,7 @@ type Monitor interface { XBlockdevChange(parent, node, child string, callback StringCallback) BlockStream(drive string, idx, blkCnt int, callback StringCallback) DriveMirror(callback StringCallback, drive, target, syncMode, format string, unmap, blockReplication bool) + DriveBackup(callback StringCallback, drive, target, syncMode, format string) BlockJobComplete(drive string, cb StringCallback) BlockReopenImage(drive, newImagePath, format string, cb StringCallback) SnapshotBlkdev(drive, newImagePath, format string, reuse bool, cb StringCallback) diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index 83dec7ce26..dc39cd67a1 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -920,6 +920,27 @@ func (m *QmpMonitor) DriveMirror(callback StringCallback, drive, target, syncMod m.Query(cmd, cb) } +func (m *QmpMonitor) DriveBackup(callback StringCallback, drive, target, syncMode, format string) { + var ( + cb = func(res *Response) { + callback(m.actionResult(res)) + } + args = map[string]interface{}{ + "device": drive, + "target": target, + "mode": "existing", + "sync": syncMode, + "format": format, + } + ) + cmd := &Command{ + Execute: "drive-backup", + Args: args, + } + + m.Query(cmd, cb) +} + func (m *QmpMonitor) BlockStream(drive string, idx, blkCnt int, callback StringCallback) { var ( speed = 5 * 100 * 1024 * 1024 // limit 500 MB/s