mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
* fix(host): use drive backup sync backup disks Signed-off-by: wanyaoqi <d3lx.yq@gmail.com> * fix(region,host): guest backup backend use fetcherfs Signed-off-by: wanyaoqi <d3lx.yq@gmail.com>
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user