diff --git a/cmd/climc/shell/compute/servers.go b/cmd/climc/shell/compute/servers.go index d7aa648bd7..c641713697 100644 --- a/cmd/climc/shell/compute/servers.go +++ b/cmd/climc/shell/compute/servers.go @@ -55,7 +55,8 @@ func init() { cmd.BatchPerform("sync", new(options.ServerIdsOptions)) cmd.Perform("switch-to-backup", new(options.ServerSwitchToBackupOptions)) cmd.BatchPerform("reconcile-backup", new(options.ServerIdsOptions)) - cmd.BatchPerform("create-backup", new(options.ServerIdsOptions)) + cmd.Perform("create-backup", new(options.ServerCreateBackupOptions)) + cmd.BatchPerform("start-backup", new(options.ServerIdsOptions)) cmd.Perform("delete-backup", new(options.ServerDeleteBackupOptions)) cmd.BatchPerform("stop", new(options.ServerStopOptions)) cmd.BatchPerform("suspend", new(options.ServerIdsOptions)) diff --git a/pkg/apis/compute/guest_const.go b/pkg/apis/compute/guest_const.go index 1bc56a71ae..67deba546b 100644 --- a/pkg/apis/compute/guest_const.go +++ b/pkg/apis/compute/guest_const.go @@ -45,7 +45,9 @@ const ( VM_DETACH_DISK = "detach_disk" VM_BACKUP_STARTING = "backup_starting" - VM_BACKUP_CREATING = compute.VM_BACKUP_CREATING + VM_BACKUP_STOPING = "backup_stopping" + VM_BACKUP_CREATING = "backup_creating" + VM_BACKUP_START_FAILED = "backup_start_failed" VM_BACKUP_CREATE_FAILED = "backup_create_fail" VM_DEPLOYING_BACKUP = "deploying_backup" VM_DEPLOYING_BACKUP_FAILED = "deploging_backup_fail" diff --git a/pkg/apis/compute/guest_metadata.go b/pkg/apis/compute/guest_metadata.go index 130d8762d6..586de770a7 100644 --- a/pkg/apis/compute/guest_metadata.go +++ b/pkg/apis/compute/guest_metadata.go @@ -15,9 +15,11 @@ package compute const ( - MIRROR_JOB = "__mirror_job_status" - MIRROR_JOB_READY = "ready" - MIRROR_JOB_FAILED = "failed" + MIRROR_JOB = "__mirror_job_status" + MIRROR_JOB_READY = "ready" + MIRROR_JOB_FAILED = "failed" + MIRROR_JOB_INPROGRESS = "inprogress" + QUORUM_CHILD_INDEX = "__quorum_child_index" DISK_CLONE_TASK_ID = "__disk_clone_task_id" diff --git a/pkg/apis/compute/guests.go b/pkg/apis/compute/guests.go index 90ca257b3b..4a5328cbfc 100644 --- a/pkg/apis/compute/guests.go +++ b/pkg/apis/compute/guests.go @@ -184,6 +184,8 @@ type ServerDetails struct { BackupHostName string `json:"backup_host_name"` // 备份主机所在宿主机状态 BackupHostStatus string `json:"backup_host_status"` + // 主备机同步状态 + BackupGuestSyncStatus string `json:"backup_guest_sync_status"` // 是否可以回收 CanRecycle bool `json:"can_recycle"` diff --git a/pkg/apis/input.go b/pkg/apis/input.go index b7c666750c..c6a2d89783 100644 --- a/pkg/apis/input.go +++ b/pkg/apis/input.go @@ -227,7 +227,8 @@ type PerformStatusInput struct { BlockJobsCount int `json:"block_jobs_count"` // 电源状态 PowerStates string `json:"power_states"` - + // is call from slave guest + IsSlave bool `json:"is_slave"` // 更改状态的原因描述 // required:false Reason string `json:"reason"` diff --git a/pkg/cloudcommon/db/opslog_const.go b/pkg/cloudcommon/db/opslog_const.go index a20e78c544..4442e10970 100644 --- a/pkg/cloudcommon/db/opslog_const.go +++ b/pkg/cloudcommon/db/opslog_const.go @@ -34,11 +34,12 @@ const ( ACT_SYNC_UPDATE = "sync_update" ACT_SYNC_CREATE = "sync_create" - ACT_START_CREATE_BACKUP = "start_create_backup" - ACT_CREATE_BACKUP = "create_backup" - ACT_CREATE_BACKUP_FAILED = "create_backup_failed" - ACT_DELETE_BACKUP = "delete_backup" - ACT_DELETE_BACKUP_FAILED = "delete_backup_failed" + ACT_START_CREATE_BACKUP = "start_create_backup" + ACT_CREATE_BACKUP = "create_backup" + ACT_CREATE_BACKUP_FAILED = "create_backup_failed" + ACT_DELETE_BACKUP = "delete_backup" + ACT_DELETE_BACKUP_FAILED = "delete_backup_failed" + ACT_UPDATE_BACKUP_GUEST_STATUS = "update_backup_guest_status" ACT_UPDATE_STATUS = "updatestatus" ACT_STARTING = "starting" diff --git a/pkg/compute/guestdrivers/kvm.go b/pkg/compute/guestdrivers/kvm.go index b657fdbc0d..83b12f2cda 100644 --- a/pkg/compute/guestdrivers/kvm.go +++ b/pkg/compute/guestdrivers/kvm.go @@ -300,10 +300,15 @@ func (self *SKVMGuestDriver) RequestStartOnHost(ctx context.Context, guest *mode config.Add(params, "params") } url := fmt.Sprintf("%s/servers/%s/start", host.ManagerUri, guest.Id) - _, _, err = httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, config, false) + _, body, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, config, false) if err != nil { return err } + if jsonutils.QueryBoolean(body, "is_running", false) { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + return body, nil + }) + } return nil } @@ -681,7 +686,7 @@ func (self *SKVMGuestDriver) RequestSyncToBackup(ctx context.Context, guest *mod body := jsonutils.NewDict() body.Add(desc, "desc") body.Set("backup_nbd_server_uri", jsonutils.NewString(guest.GetMetadata(ctx, "backup_nbd_server_uri", task.GetUserCred()))) - url := fmt.Sprintf("%s/servers/%s/drive-mirror", host.ManagerUri, guest.Id) + url := fmt.Sprintf("%s/servers/%s/block-replication", host.ManagerUri, guest.Id) header := self.getTaskRequestHeader(task) _, _, err = httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, header, body, false) if err != nil { diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 8a9ba10a92..86cff00853 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -2835,14 +2835,40 @@ func (self *SGuest) SetPowerStates(powerStates string) error { return errors.Wrap(err, "Update power states") } +func (self *SGuest) SetBackupGuestStatus(userCred mcclient.TokenCredential, status string, reason string) error { + if self.BackupGuestStatus == status { + return nil + } + oldStatus := self.BackupGuestStatus + _, err := db.Update(self, func() error { + self.BackupGuestStatus = status + return nil + }) + if err != nil { + return errors.Wrap(err, "Update backup guest status") + } + if userCred != nil { + notes := fmt.Sprintf("%s=>%s", oldStatus, status) + if len(reason) > 0 { + notes = fmt.Sprintf("%s: %s", notes, reason) + } + db.OpsLog.LogEvent(self, db.ACT_UPDATE_BACKUP_GUEST_STATUS, notes, userCred) + logclient.AddSimpleActionLog(self, logclient.ACT_UPDATE_BACKUP_GUEST_STATUS, notes, userCred, true) + } + return nil +} + func (self *SGuest) PerformStatus(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input apis.PerformStatusInput) (jsonutils.JSONObject, error) { + if input.IsSlave { // perform status called from slave guest + return nil, self.SetBackupGuestStatus(userCred, input.Status, input.Reason) + } if input.PowerStates != "" { if err := self.SetPowerStates(input.PowerStates); err != nil { return nil, errors.Wrap(err, "set power states") } } - preStatus := self.Status + preStatus := self.Status if len(self.BackupHostId) == 0 && input.Status == api.VM_RUNNING && input.BlockJobsCount > 0 { input.Status = api.VM_BLOCK_STREAM } @@ -2851,9 +2877,23 @@ func (self *SGuest) PerformStatus(ctx context.Context, userCred mcclient.TokenCr return nil, errors.Wrap(err, "SVirtualResourceBase.PerformStatus") } - 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 { + 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 err := self.ResetGuestQuorumChildIndex(ctx, userCred); err != nil { + return nil, errors.Wrap(err, "reset guest quorum child index") + } + } + } + + if ispId := self.GetMetadata(ctx, api.BASE_INSTANCE_SNAPSHOT_ID, userCred); len(ispId) > 0 { ispM, err := InstanceSnapshotManager.FetchById(ispId) if err == nil { isp := ispM.(*SInstanceSnapshot) @@ -3249,27 +3289,31 @@ func (self *SGuest) SwitchToBackup(userCred mcclient.TokenCredential) error { return nil } -func (self *SGuest) PerformSwitchToBackup(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { +func (self *SGuest) PerformSwitchToBackup( + ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject, +) (jsonutils.JSONObject, error) { if self.Status == api.VM_BLOCK_STREAM { return nil, httperrors.NewBadRequestError("Cannot swith to backup when guest in status %s", self.Status) } if len(self.BackupHostId) == 0 { return nil, httperrors.NewBadRequestError("Guest no backup host") } + backupHost := HostManager.FetchHostById(self.BackupHostId) + if backupHost.HostStatus != api.HOST_ONLINE { + return nil, httperrors.NewBadRequestError("Can't switch to backup host on host status %s", backupHost.HostStatus) + } - mirrorJobStatus := self.GetMetadata(ctx, api.MIRROR_JOB, userCred) - if mirrorJobStatus != api.MIRROR_JOB_READY { + if !self.IsGuestBackupMirrorJobReady(ctx, userCred) { return nil, httperrors.NewBadRequestError("Guest can't switch to backup, mirror job not ready") } + if !utils.IsInStringArray(self.BackupGuestStatus, []string{api.VM_RUNNING, api.VM_READY, api.VM_UNKNOWN}) { + return nil, httperrors.NewInvalidStatusError("Guest can't switch to backup with backup status %s", self.BackupGuestStatus) + } oldStatus := self.Status - deleteBackup := jsonutils.QueryBoolean(data, "delete_backup", false) - purgeBackup := jsonutils.QueryBoolean(data, "purge_backup", false) - taskData := jsonutils.NewDict() taskData.Set("old_status", jsonutils.NewString(oldStatus)) - taskData.Set("delete_backup", jsonutils.NewBool(deleteBackup)) - taskData.Set("purge_backup", jsonutils.NewBool(purgeBackup)) + taskData.Set("auto_start", jsonutils.NewBool(jsonutils.QueryBoolean(data, "auto_start", false))) if task, err := taskman.TaskManager.NewTask(ctx, "GuestSwitchToBackupTask", self, userCred, taskData, "", "", nil); err != nil { log.Errorln(err) return nil, err @@ -3365,7 +3409,9 @@ func (manager *SGuestManager) PerformBatchSetUserMetadata(ctx context.Context, u func (self *SGuest) PerformBlockStreamFailed(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { if len(self.BackupHostId) > 0 { - self.SetMetadata(ctx, api.MIRROR_JOB, api.MIRROR_JOB_FAILED, userCred) + if err := self.SetGuestBackupMirrorJobFailed(ctx, userCred); err != nil { + return nil, errors.Wrap(err, "set guest backup mirror job failed") + } } if self.Status == api.VM_BLOCK_STREAM || self.Status == api.VM_RUNNING { reason, _ := data.GetString("reason") @@ -3523,6 +3569,9 @@ func (self *SGuest) PerformCreateBackup( if len(self.BackupHostId) > 0 { return nil, httperrors.NewBadRequestError("Already have backup server") } + if self.Status != api.VM_READY { + return nil, httperrors.NewBadRequestError("Can't create backup in guest status %s", self.Status) + } if !self.guestDisksStorageTypeIsLocal() { return nil, httperrors.NewBadRequestError("Cannot create backup with shared storage") } @@ -3559,6 +3608,7 @@ func (self *SGuest) StartGuestCreateBackupTask( params := data.(*jsonutils.JSONDict) params.Set("guest_status", jsonutils.NewString(self.Status)) + self.SetStatus(userCred, api.VM_BACKUP_CREATING, "") task, err := taskman.TaskManager.NewTask(ctx, "GuestCreateBackupTask", self, userCred, params, parentTaskId, "", &req) if err != nil { quotas.CancelPendingUsage(ctx, userCred, &req, &req, false) @@ -3620,38 +3670,19 @@ func (self *SGuest) StartCreateBackup(ctx context.Context, userCred mcclient.Tok return nil } -func (self *SGuest) PerformReconcileBackup(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { - switchBackup := self.GetMetadata(ctx, "switch_backup", userCred) - createBackup := self.GetMetadata(ctx, "create_backup", userCred) - if len(switchBackup) == 0 && len(createBackup) == 0 { - return nil, httperrors.NewBadRequestError("guest doesn't need reconcile backup") +func (self *SGuest) PerformStartBackup( + ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject, +) (jsonutils.JSONObject, error) { + if !self.HasBackupGuest() { + return nil, httperrors.NewBadRequestError("guest has no backup guest") } - if len(switchBackup) > 0 { - data := jsonutils.NewDict() - data.Set("purge_backup", jsonutils.JSONTrue) - return self.PerformSwitchToBackup(ctx, userCred, nil, data) - } else { - return nil, self.StartReconcileBackup(ctx, userCred) + if host := HostManager.FetchHostById(self.BackupHostId); host.HostStatus != api.HOST_ONLINE { + return nil, httperrors.NewBadRequestError("can't start backup guest on host status %s", host.HostStatus) } -} - -func (self *SGuest) StartReconcileBackup(ctx context.Context, userCred mcclient.TokenCredential) error { - if len(self.BackupHostId) > 0 { - data := jsonutils.NewDict() - data.Set("purge", jsonutils.JSONTrue) - data.Set("create", jsonutils.JSONTrue) - _, err := self.PerformDeleteBackup(ctx, userCred, nil, data) - if err != nil { - return err - } - } else { - params := jsonutils.NewDict() - params.Set("reconcile_backup", jsonutils.JSONTrue) - if _, err := self.StartGuestCreateBackupTask(ctx, userCred, "", params); err != nil { - return err - } + if self.Status != api.VM_RUNNING || self.BackupGuestStatus == api.VM_RUNNING { + return nil, httperrors.NewBadRequestError("can't start backup guest on backup guest status %s", self.BackupGuestStatus) } - return nil + return nil, self.GuestStartAndSyncToBackup(ctx, userCred, "", self.Status) } func (self *SGuest) PerformSetExtraOption(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, input api.ServerSetExtraOptionInput) (jsonutils.JSONObject, error) { diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 6741277171..98e22ee345 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -135,7 +135,8 @@ type SGuest struct { KeypairId string `width:"36" charset:"ascii" nullable:"true" list:"user" create:"optional"` // 备份机所在宿主机Id - BackupHostId string `width:"36" charset:"ascii" nullable:"true" list:"user" get:"user"` + BackupHostId string `width:"36" charset:"ascii" nullable:"true" list:"user" get:"user"` + BackupGuestStatus string `width:"36" charset:"ascii" nullable:"false" default:"init" list:"user" create:"optional" json:"backup_guest_status"` // 迁移或克隆的速度 ProgressMbps float64 `nullable:"false" default:"0" list:"user" create:"optional" update:"user" log:"skip"` @@ -2264,6 +2265,7 @@ func (self *SGuest) moreExtraInfo( if len(fields) == 0 || fields.Contains("backup_host_status") { out.BackupHostStatus = backupHost.HostStatus } + out.BackupGuestSyncStatus = self.GetGuestBackupMirrorJobStatus(ctx, userCred) } if len(fields) == 0 || fields.Contains("can_recycle") { @@ -5371,115 +5373,6 @@ func (manager *SGuestManager) DeleteExpiredPostpaidServers(ctx context.Context, } } -func (self *SGuestManager) ReconcileBackupGuests(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { - self.switchBackupGuests(ctx, userCred) - self.createBackupGuests(ctx, userCred) -} - -func (self *SGuestManager) switchBackupGuests(ctx context.Context, userCred mcclient.TokenCredential) { - q := self.Query() - metaDataQuery := db.Metadata.Query().Startswith("id", "server::"). - Equals("key", "switch_backup").IsNotEmpty("value").GroupBy("id") - metaDataQuery.AppendField(sqlchemy.SubStr("guest_id", metaDataQuery.Field("id"), len("server::")+1, 0)) - subQ := metaDataQuery.SubQuery() - q = q.Join(subQ, sqlchemy.Equals(q.Field("id"), subQ.Field("guest_id"))) - - guests := make([]SGuest, 0) - err := db.FetchModelObjects(GuestManager, q, &guests) - if err != nil { - log.Errorf("ReconcileBackupGuests failed fetch guests %s", err) - return - } - log.Debugf("Guests count %d need reconcile with switch backup", len(guests)) - for i := 0; i < len(guests); i++ { - val := guests[i].GetMetadataJson(ctx, "switch_backup", userCred) - t, err := val.GetTime() - if err != nil { - log.Errorf("failed get time from metadata switch_backup %s", err) - continue - } - if time.Now().After(t) { - data := jsonutils.NewDict() - data.Set("purge_backup", jsonutils.JSONTrue) - _, err := guests[i].PerformSwitchToBackup(ctx, userCred, nil, data) - if err != nil { - db.OpsLog.LogEvent( - &guests[i], db.ACT_SWITCH_FAILED, fmt.Sprintf("switchBackupGuests on reconcile_backup: %s", err), userCred, - ) - logclient.AddSimpleActionLog( - &guests[i], logclient.ACT_SWITCH_TO_BACKUP, - fmt.Sprintf("switchBackupGuests on reconcile_backup: %s", err), userCred, false, - ) - } - } - } -} - -func (self *SGuestManager) createBackupGuests(ctx context.Context, userCred mcclient.TokenCredential) { - q := self.Query() - metaDataQuery := db.Metadata.Query().Startswith("id", "server::"). - Equals("key", "create_backup").IsNotEmpty("value").GroupBy("id") - metaDataQuery.AppendField(sqlchemy.SubStr("guest_id", metaDataQuery.Field("id"), len("server::")+1, 0)) - subQ := metaDataQuery.SubQuery() - q = q.Join(subQ, sqlchemy.Equals(q.Field("id"), subQ.Field("guest_id"))) - - guests := make([]SGuest, 0) - err := db.FetchModelObjects(GuestManager, q, &guests) - if err != nil { - log.Errorf("ReconcileBackupGuests failed fetch guests %s", err) - return - } - log.Infof("Guests count %d need reconcile with create bakcup", len(guests)) - for i := 0; i < len(guests); i++ { - val := guests[i].GetMetadataJson(ctx, "create_backup", userCred) - t, err := val.GetTime() - if err != nil { - log.Errorf("failed get time from metadata create_backup %s", err) - continue - } - if time.Now().After(t) { - if len(guests[i].BackupHostId) > 0 { - data := jsonutils.NewDict() - data.Set("purge", jsonutils.JSONTrue) - data.Set("create", jsonutils.JSONTrue) - _, err := guests[i].PerformDeleteBackup(ctx, userCred, nil, data) - if err != nil { - db.OpsLog.LogEvent( - &guests[i], db.ACT_DELETE_BACKUP_FAILED, - fmt.Sprintf("PerformDeleteBackup on ReconcileBackupGuests: %s", err), userCred, - ) - logclient.AddSimpleActionLog( - &guests[i], logclient.ACT_DELETE_BACKUP, - fmt.Sprintf("PerformDeleteBackup on ReconcileBackupGuests: %s", err), userCred, false, - ) - } - } else { - params := jsonutils.NewDict() - params.Set("reconcile_backup", jsonutils.JSONTrue) - if _, err := guests[i].StartGuestCreateBackupTask(ctx, userCred, "", params); err != nil { - db.OpsLog.LogEvent( - &guests[i], db.ACT_CREATE_BACKUP_FAILED, - fmt.Sprintf("StartGuestCreateBackupTask on ReconcileBackupGuests: %s", err), userCred, - ) - logclient.AddSimpleActionLog( - &guests[i], logclient.ACT_CREATE_BACKUP, - fmt.Sprintf("StartGuestCreateBackupTask on ReconcileBackupGuests: %s", err), userCred, false, - ) - } - } - } - } -} - -func (self *SGuest) isInReconcile(userCred mcclient.TokenCredential) bool { - switchBackup := self.GetMetadata(context.Background(), "switch_backup", userCred) - createBackup := self.GetMetadata(context.Background(), "create_backup", userCred) - if len(switchBackup) > 0 || len(createBackup) > 0 { - return true - } - return false -} - func (self *SGuest) IsEipAssociable() error { if !utils.IsInStringArray(self.Status, []string{api.VM_READY, api.VM_RUNNING}) { return errors.Wrapf(httperrors.ErrInvalidStatus, "cannot associate eip in status %s", self.Status) @@ -6411,3 +6304,43 @@ func (guest *SGuest) inferPowerStates() { } } } + +func (guest *SGuest) HasBackupGuest() bool { + return guest.BackupHostId != "" +} + +func (guest *SGuest) SetGuestBackupMirrorJobInProgress(ctx context.Context, userCred mcclient.TokenCredential) error { + return guest.SetMetadata(ctx, api.MIRROR_JOB, api.MIRROR_JOB_INPROGRESS, userCred) +} + +func (guest *SGuest) SetGuestBackupMirrorJobNotReady(ctx context.Context, userCred mcclient.TokenCredential) error { + return guest.SetMetadata(ctx, api.MIRROR_JOB, "", userCred) +} + +func (guest *SGuest) TrySetGuestBackupMirrorJobReady(ctx context.Context, userCred mcclient.TokenCredential) error { + if guest.IsGuestBackupMirrorJobFailed(ctx, userCred) { + // can't update guest backup mirror job status from failed to ready + return nil + } + return guest.SetMetadata(ctx, api.MIRROR_JOB, api.MIRROR_JOB_READY, userCred) +} + +func (guest *SGuest) SetGuestBackupMirrorJobFailed(ctx context.Context, userCred mcclient.TokenCredential) error { + return guest.SetMetadata(ctx, api.MIRROR_JOB, api.MIRROR_JOB_FAILED, userCred) +} + +func (guest *SGuest) IsGuestBackupMirrorJobFailed(ctx context.Context, userCred mcclient.TokenCredential) bool { + return guest.GetMetadata(ctx, api.MIRROR_JOB, userCred) == api.MIRROR_JOB_FAILED +} + +func (guest *SGuest) IsGuestBackupMirrorJobReady(ctx context.Context, userCred mcclient.TokenCredential) bool { + return guest.GetMetadata(ctx, api.MIRROR_JOB, userCred) == api.MIRROR_JOB_READY +} + +func (guest *SGuest) GetGuestBackupMirrorJobStatus(ctx context.Context, userCred mcclient.TokenCredential) string { + return guest.GetMetadata(ctx, api.MIRROR_JOB, userCred) +} + +func (guest *SGuest) ResetGuestQuorumChildIndex(ctx context.Context, userCred mcclient.TokenCredential) error { + return guest.SetMetadata(ctx, api.QUORUM_CHILD_INDEX, "", userCred) +} diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index f8a85c6fb1..ca815c82cc 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -5519,6 +5519,10 @@ func (self *SHost) MarkGuestUnknown(userCred mcclient.TokenCredential) { for _, guest := range guests { guest.SetStatus(userCred, api.VM_UNKNOWN, "host offline") } + guests2 := self.GetGuestsBackupOnThisHost() + for _, guest := range guests2 { + guest.SetBackupGuestStatus(userCred, api.VM_UNKNOWN, "host offline") + } } func (manager *SHostManager) PingDetectionTask(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { @@ -5724,12 +5728,7 @@ func (host *SHost) OnHostDown(ctx context.Context, userCred mcclient.TokenCreden func (host *SHost) switchWithBackup(ctx context.Context, userCred mcclient.TokenCredential) { guests := host.GetGuestsMasterOnThisHost() for i := 0; i < len(guests); i++ { - if guests[i].isInReconcile(userCred) { - log.Warningf("guest %s is in reconcile", guests[i].GetName()) - continue - } data := jsonutils.NewDict() - data.Set("purge_backup", jsonutils.JSONTrue) _, err := guests[i].PerformSwitchToBackup(ctx, userCred, nil, data) if err != nil { db.OpsLog.LogEvent( @@ -5739,17 +5738,11 @@ func (host *SHost) switchWithBackup(ctx context.Context, userCred mcclient.Token &guests[i], logclient.ACT_SWITCH_TO_BACKUP, fmt.Sprintf("PerformSwitchToBackup on host down: %s", err), userCred, false, ) - } else { - guests[i].SetMetadata(ctx, "origin_status", guests[i].Status, userCred) } } guests2 := host.GetGuestsBackupOnThisHost() for i := 0; i < len(guests2); i++ { - if guests2[i].isInReconcile(userCred) { - log.Warningf("guest %s is in reconcile", guests2[i].GetName()) - continue - } data := jsonutils.NewDict() data.Set("purge", jsonutils.JSONTrue) data.Set("create", jsonutils.JSONTrue) diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index d1b7975666..889e1bc94d 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -156,11 +156,6 @@ func StartService() { cron.AddJobAtIntervalsWithStartRun("CalculateDomainQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.DomainQuotaManager.CalculateQuotaUsages, true) cron.AddJobAtIntervalsWithStartRun("CalculateInfrasQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.InfrasQuotaManager.CalculateQuotaUsages, true) cron.AddJobAtIntervalsWithStartRun("AutoSyncCloudaccountStatusTask", time.Duration(opts.CloudAutoSyncIntervalSeconds)*time.Second, models.CloudaccountManager.AutoSyncCloudaccountStatusTask, true) - - if opts.AutoReconcileBackupServers { - cron.AddJobAtIntervalsWithStartRun("ReconcileBackupGuests", time.Duration(opts.ReconcileGuestBackupIntervalSeconds)*time.Second, models.GuestManager.ReconcileBackupGuests, true) - } - cron.AddJobAtIntervalsWithStartRun("SyncCapacityUsedForEsxiStorage", time.Duration(opts.SyncStorageCapacityUsedIntervalMinutes)*time.Minute, models.StorageManager.SyncCapacityUsedForEsxiStorage, true) cron.AddJobAtIntervalsWithStartRun("AutoSyncExtDiskSnapshot", time.Duration(opts.SyncExtDiskSnapshotIntervalMinutes)*time.Minute, models.DiskManager.AutoSyncExtDiskSnapshot, true) diff --git a/pkg/compute/tasks/guest_backup_tasks.go b/pkg/compute/tasks/guest_backup_tasks.go index 5c87f5dc31..a8ae945773 100644 --- a/pkg/compute/tasks/guest_backup_tasks.go +++ b/pkg/compute/tasks/guest_backup_tasks.go @@ -17,10 +17,8 @@ package tasks import ( "context" "fmt" - "time" "yunion.io/x/jsonutils" - "yunion.io/x/log" "yunion.io/x/pkg/utils" api "yunion.io/x/onecloud/pkg/apis/compute" @@ -88,116 +86,45 @@ 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) - - self.SetStage("OnSwitched", nil) - if jsonutils.QueryBoolean(self.Params, "purge_backup", false) { - guest.StartGuestDeleteOnHostTask(ctx, self.UserCred, guest.BackupHostId, true, self.GetTaskId()) - } else if jsonutils.QueryBoolean(self.Params, "delete_backup", false) { - guest.StartGuestDeleteOnHostTask(ctx, self.UserCred, guest.BackupHostId, false, self.GetTaskId()) + oldStatus, _ := self.Params.GetString("old_status") + autoStart := jsonutils.QueryBoolean(self.Params, "auto_start", false) || + utils.IsInStringArray(oldStatus, api.VM_RUNNING_STATUS) + if autoStart { + self.SetStage("OnGuestStartCompleted", nil) + if err := guest.StartGueststartTask(ctx, self.UserCred, nil, self.GetId()); err != nil { + self.OnGuestStartCompletedFailed(ctx, guest, + jsonutils.NewString(fmt.Sprintf("start guest start task: %s", err))) + } } else { - self.OnSwitched(ctx, guest, nil) + self.OnComplete(ctx, guest, nil) } } -func (self *GuestSwitchToBackupTask) OnNewMasterStarted(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - guest.RemoveMetadata(ctx, "origin_status", self.UserCred) - self.OnComplete(ctx, guest, nil) -} - func (self *GuestSwitchToBackupTask) OnFail(ctx context.Context, guest *models.SGuest, reason jsonutils.JSONObject) { guest.SetStatus(self.UserCred, api.VM_SWITCH_TO_BACKUP_FAILED, reason.String()) db.OpsLog.LogEvent(guest, db.ACT_SWITCH_FAILED, reason, self.UserCred) logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_SWITCH_TO_BACKUP, reason, self.UserCred, false) - self.SetSwitchFiledGuestMetadata(ctx, guest) self.SetStageFailed(ctx, reason) } -func (self *GuestSwitchToBackupTask) SetSwitchFiledGuestMetadata(ctx context.Context, guest *models.SGuest) { - if res := guest.GetMetadata(ctx, "switch_backup", self.UserCred); len(res) == 0 { - guest.SetMetadata( - ctx, "switch_backup", jsonutils.NewTimeString(time.Now().Add(time.Minute*1).UTC()), self.UserCred) - guest.SetMetadata(ctx, "switch_backup_count", jsonutils.NewInt(1), self.UserCred) - } else { - count := guest.GetMetadataJson(ctx, "switch_backup_count", self.UserCred) - cnt, _ := count.Int() - cnt += 1 - dur := cnt - if dur > 15 { - dur = 15 - } - - guest.SetMetadata(ctx, "switch_backup", - jsonutils.NewTimeString(time.Now().Add(time.Minute*time.Duration(dur)).UTC()), self.UserCred) - guest.SetMetadata(ctx, "switch_backup_count", jsonutils.NewInt(cnt), self.UserCred) - } +func (self *GuestSwitchToBackupTask) OnGuestStartCompleted(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.SetStageComplete(ctx, nil) } -func (self *GuestSwitchToBackupTask) CleanGuestMetadata(ctx context.Context, guest *models.SGuest) { - if res := guest.GetMetadata(ctx, "switch_backup", self.UserCred); len(res) > 0 { - guest.RemoveMetadata(ctx, "switch_backup", self.UserCred) - guest.RemoveMetadata(ctx, "switch_backup_count", self.UserCred) - } +func (self *GuestSwitchToBackupTask) OnGuestStartCompletedFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.SetStageFailed(ctx, data) } func (self *GuestSwitchToBackupTask) OnComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { self.SetStageComplete(ctx, nil) } -func (self *GuestSwitchToBackupTask) OnSwitched(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - if err := guest.SetMetadata(ctx, api.MIRROR_JOB, "", self.UserCred); err != nil { - self.OnSwitchedFailed(ctx, guest, jsonutils.NewString("guest set metadata failed")) - return - } - // switched bakcup success - self.CleanGuestMetadata(ctx, guest) - - // add a new backup server - self.SetStage("OnCreatedBackup", nil) - if jsonutils.QueryBoolean(self.Params, "purge_backup", false) || - jsonutils.QueryBoolean(self.Params, "delete_backup", false) { - params := jsonutils.NewDict() - params.Set("reconcile_backup", jsonutils.JSONTrue) - if _, err := guest.StartGuestCreateBackupTask(ctx, self.UserCred, self.Id, params); err != nil { - log.Errorf("guest start create backup failed %s", err) - self.failedStartCreateBackupTask(ctx, guest) - } - } else { - self.OnCreatedBackup(ctx, guest, nil) - } -} - -func (self *GuestSwitchToBackupTask) failedStartCreateBackupTask(ctx context.Context, guest *models.SGuest) { - guest.SetMetadata(ctx, "create_backup", jsonutils.NewTimeString(time.Now().Add(time.Minute*1).UTC()), self.UserCred) - guest.SetMetadata(ctx, "create_backup_count", jsonutils.NewInt(1), self.UserCred) - self.OnCreatedBackup(ctx, guest, nil) -} - -func (self *GuestSwitchToBackupTask) OnCreatedBackup(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - oldStatus, _ := self.Params.GetString("old_status") - originStatus := guest.GetMetadata(ctx, "origin_status", self.UserCred) - if (!utils.IsInStringArray(guest.Status, api.VM_RUNNING_STATUS) && utils.IsInStringArray(oldStatus, api.VM_RUNNING_STATUS)) || - utils.IsInStringArray(originStatus, api.VM_RUNNING_STATUS) { - self.SetStage("OnNewMastqerStarted", nil) - guest.StartGueststartTask(ctx, self.UserCred, nil, self.GetTaskId()) - } else { - self.OnComplete(ctx, guest, nil) - } - if len(originStatus) > 0 { - guest.RemoveMetadata(ctx, "origin_status", self.UserCred) - } -} - -func (self *GuestSwitchToBackupTask) OnCreatedBackupFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - log.Errorf("failed create backup %s", data) - self.OnCreatedBackup(ctx, guest, data) -} - -func (self *GuestSwitchToBackupTask) OnSwitchedFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - self.OnFail(ctx, guest, data) -} - /********************* GuestStartAndSyncToBackupTask *********************/ type GuestStartAndSyncToBackupTask struct { @@ -206,6 +133,7 @@ type GuestStartAndSyncToBackupTask struct { func (self *GuestStartAndSyncToBackupTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { guest := obj.(*models.SGuest) + guest.SetStatus(self.UserCred, api.VM_BACKUP_STARTING, "GuestStartAndSyncToBackupTask") self.SetStage("OnCheckTemplete", nil) self.checkTemplete(ctx, guest) } @@ -233,6 +161,7 @@ func (self *GuestStartAndSyncToBackupTask) OnCheckTemplete(ctx context.Context, } func (self *GuestStartAndSyncToBackupTask) OnStartBackupGuest(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + guest.SetBackupGuestStatus(self.UserCred, api.VM_RUNNING, "on start backup guest") nbdServerPort, err := data.Int("nbd_server_port") if err != nil { self.SetStageFailed(ctx, jsonutils.NewString("Start Backup Guest Missing Nbd Port")) @@ -261,18 +190,19 @@ func (self *GuestStartAndSyncToBackupTask) OnStartBackupGuest(ctx context.Contex } func (self *GuestStartAndSyncToBackupTask) OnStartBackupGuestFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - guest.SetMetadata(ctx, api.MIRROR_JOB, api.MIRROR_JOB_FAILED, self.UserCred) + guest.SetGuestBackupMirrorJobFailed(ctx, self.UserCred) db.OpsLog.LogEvent(guest, db.ACT_BACKUP_START_FAILED, data.String(), self.UserCred) self.SetStageFailed(ctx, data) } func (self *GuestStartAndSyncToBackupTask) OnRequestSyncToBackup(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - guest.SetMetadata(ctx, api.MIRROR_JOB, "", self.UserCred) + guest.SetGuestBackupMirrorJobInProgress(ctx, self.UserCred) guest.SetStatus(self.UserCred, api.VM_BLOCK_STREAM, "OnSyncToBackup") self.SetStageComplete(ctx, nil) } func (self *GuestStartAndSyncToBackupTask) OnRequestSyncToBackupFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + guest.SetGuestBackupMirrorJobFailed(ctx, self.UserCred) guest.SetStatus(self.UserCred, api.VM_BLOCK_STREAM_FAIL, "OnSyncToBackup") self.SetStageFailed(ctx, data) } @@ -287,7 +217,6 @@ func (self *GuestCreateBackupTask) OnInit(ctx context.Context, obj db.IStandalon func (self *GuestCreateBackupTask) OnStartSchedule(obj IScheduleModel) { guest := obj.(*models.SGuest) - guest.SetStatus(self.UserCred, api.VM_BACKUP_CREATING, "") db.OpsLog.LogEvent(guest, db.ACT_START_CREATE_BACKUP, "", self.UserCred) } @@ -389,9 +318,12 @@ func (self *GuestCreateBackupTask) OnCreateBackupFailed(ctx context.Context, gue } func (self *GuestCreateBackupTask) OnCreateBackup(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + guest.SetBackupGuestStatus(self.UserCred, api.VM_READY, "on create backup") guestStatus, _ := self.Params.GetString("guest_status") if utils.IsInStringArray(guestStatus, api.VM_RUNNING_STATUS) { self.OnGuestStart(ctx, guest, guestStatus) + } else if jsonutils.QueryBoolean(self.Params, "auto_start", false) { + self.RequestStartGuest(ctx, guest) } else { self.TaskCompleted(ctx, guest, "") } @@ -413,14 +345,23 @@ func (self *GuestCreateBackupTask) OnSyncToBackupFailed(ctx context.Context, gue self.TaskFailed(ctx, guest, data) } +func (self *GuestCreateBackupTask) RequestStartGuest(ctx context.Context, guest *models.SGuest) { + self.SetStage("OnGuestStartCompleted", nil) + guest.StartGueststartTask(ctx, self.UserCred, nil, self.GetId()) +} + +func (self *GuestCreateBackupTask) OnGuestStartCompleted(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.TaskCompleted(ctx, guest, "") +} + +func (self *GuestCreateBackupTask) OnGuestStartCompletedFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + db.OpsLog.LogEvent(guest, db.ACT_CREATE_BACKUP_FAILED, data.String(), self.UserCred) + logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_CREATE_BACKUP, data.String(), self.UserCred, false) + self.SetStageFailed(ctx, data) +} + func (self *GuestCreateBackupTask) TaskCompleted(ctx context.Context, guest *models.SGuest, reason string) { db.OpsLog.LogEvent(guest, db.ACT_CREATE_BACKUP, reason, self.UserCred) - - if res := guest.GetMetadata(ctx, "create_backup", self.UserCred); len(res) > 0 { - guest.RemoveMetadata(ctx, "create_backup", self.UserCred) - guest.RemoveMetadata(ctx, "create_backup_count", self.UserCred) - } - logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_CREATE_BACKUP, reason, self.UserCred, true) self.SetStageComplete(ctx, nil) guest.StartSyncstatus(ctx, self.UserCred, "") @@ -428,25 +369,6 @@ func (self *GuestCreateBackupTask) TaskCompleted(ctx context.Context, guest *mod func (self *GuestCreateBackupTask) TaskFailed(ctx context.Context, guest *models.SGuest, reason jsonutils.JSONObject) { guest.SetStatus(self.UserCred, api.VM_BACKUP_CREATE_FAILED, reason.String()) - if jsonutils.QueryBoolean(self.Params, "reconcile_backup", false) { - if res := guest.GetMetadata(ctx, "create_backup", self.UserCred); len(res) == 0 { - guest.SetMetadata( - ctx, "create_backup", jsonutils.NewTimeString(time.Now().Add(time.Minute*1).UTC()), self.UserCred) - guest.SetMetadata(ctx, "create_backup_count", jsonutils.NewInt(1), self.UserCred) - } else { - count := guest.GetMetadataJson(ctx, "create_backup_count", self.UserCred) - cnt, _ := count.Int() - cnt += 1 - dur := cnt - if dur > 15 { - dur = 15 - } - - guest.SetMetadata(ctx, "create_backup", - jsonutils.NewTimeString(time.Now().Add(time.Minute*time.Duration(dur)).UTC()), self.UserCred) - guest.SetMetadata(ctx, "create_backup_count", jsonutils.NewInt(cnt), self.UserCred) - } - } db.OpsLog.LogEvent(guest, db.ACT_CREATE_BACKUP_FAILED, reason, self.UserCred) logclient.AddActionLogWithContext(ctx, guest, logclient.ACT_CREATE_BACKUP, reason, self.UserCred, false) self.SetStageFailed(ctx, reason) diff --git a/pkg/compute/tasks/guest_delete_backup_task.go b/pkg/compute/tasks/guest_delete_backup_task.go index 022cf4c1a3..d50a3f13a3 100644 --- a/pkg/compute/tasks/guest_delete_backup_task.go +++ b/pkg/compute/tasks/guest_delete_backup_task.go @@ -57,8 +57,8 @@ func (self *GuestDeleteBackupTask) OnInit(ctx context.Context, obj db.IStandalon return } - self.SetStage("OnCancelBlockJobs", nil) - url := fmt.Sprintf("%s/servers/%s/cancel-block-jobs", host.ManagerUri, guest.Id) + self.SetStage("OnCancelBlockReplication", nil) + url := fmt.Sprintf("%s/servers/%s/cancel-block-replication", host.ManagerUri, guest.Id) _, _, err := httputils.JSONRequest(httputils.GetDefaultClient(), ctx, "POST", url, self.GetTaskRequestHeader(), nil, false) if err != nil { @@ -67,7 +67,7 @@ func (self *GuestDeleteBackupTask) OnInit(ctx context.Context, obj db.IStandalon } } -func (self *GuestDeleteBackupTask) OnCancelBlockJobs(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { +func (self *GuestDeleteBackupTask) OnCancelBlockReplication(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { self.StartDeleteBackupOnHost(ctx, guest) } @@ -87,17 +87,17 @@ func (self *GuestDeleteBackupTask) StartDeleteBackupOnHost(ctx context.Context, } } -func (self *GuestDeleteBackupTask) OnCancelBlockJobsFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { +func (self *GuestDeleteBackupTask) OnCancelBlockReplicationFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { self.OnFail(ctx, guest, data) } func (self *GuestDeleteBackupTask) OnDeleteOnHost(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + guest.SetGuestBackupMirrorJobNotReady(ctx, self.UserCred) if jsonutils.QueryBoolean(self.Params, "create", false) { self.OnDeleteBackupComplete(ctx, guest, data) self.SetStage("OnCreateNewBackup", nil) params := jsonutils.NewDict() - params.Set("reconcile_backup", jsonutils.JSONTrue) _, err := guest.StartGuestCreateBackupTask(ctx, self.UserCred, self.GetId(), params) if err != nil { self.onCreateNewBackupFailed(ctx, guest, jsonutils.NewString(err.Error())) @@ -133,10 +133,6 @@ func (self *GuestDeleteBackupTask) OnDeleteBackupComplete(ctx context.Context, g } func (self *GuestDeleteBackupTask) TaskComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - guest.RemoveMetadata(ctx, "switch_backup", self.UserCred) - guest.RemoveMetadata(ctx, "switch_backup_count", self.UserCred) - guest.RemoveMetadata(ctx, "create_backup", self.UserCred) - guest.RemoveMetadata(ctx, "create_backup_count", self.UserCred) self.OnDeleteBackupComplete(ctx, guest, data) self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/guest_delete_on_host_task.go b/pkg/compute/tasks/guest_delete_on_host_task.go index 0d0a4dd50b..3b05da184f 100644 --- a/pkg/compute/tasks/guest_delete_on_host_task.go +++ b/pkg/compute/tasks/guest_delete_on_host_task.go @@ -20,6 +20,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/compute/models" @@ -76,6 +77,7 @@ func (self *GuestDeleteOnHostTask) OnUnDeployGuest(ctx context.Context, guest *m if guest.BackupHostId == hostId { _, err := db.Update(guest, func() error { guest.BackupHostId = "" + guest.BackupGuestStatus = compute.VM_INIT return nil }) if err != nil { diff --git a/pkg/compute/tasks/ha_guest_deploy_task.go b/pkg/compute/tasks/ha_guest_deploy_task.go index 3d6d7e3c5b..4f19330c7b 100644 --- a/pkg/compute/tasks/ha_guest_deploy_task.go +++ b/pkg/compute/tasks/ha_guest_deploy_task.go @@ -39,7 +39,12 @@ type HAGuestDeployTask struct { func (self *HAGuestDeployTask) OnDeployWaitServerStop( ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, ) { - self.DeployBackup(ctx, guest, nil) + host := models.HostManager.FetchHostById(guest.BackupHostId) + if host.HostStatus != api.HOST_ONLINE { + self.GuestDeployTask.OnDeployWaitServerStop(ctx, guest, data) + } else { + self.DeployBackup(ctx, guest, nil) + } } func (self *HAGuestDeployTask) DeployBackup(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { @@ -57,6 +62,7 @@ func (self *HAGuestDeployTask) DeployBackup(ctx context.Context, guest *models.S func (self *HAGuestDeployTask) OnDeploySlaveGuestComplete( ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, ) { + guest.SetGuestBackupMirrorJobNotReady(ctx, self.UserCred) host, _ := guest.GetHost() self.SetStage("OnDeployGuestComplete", nil) self.DeployOnHost(ctx, guest, host) @@ -65,6 +71,7 @@ func (self *HAGuestDeployTask) OnDeploySlaveGuestComplete( func (self *HAGuestDeployTask) OnDeploySlaveGuestCompleteFailed( ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, ) { + guest.SetGuestBackupMirrorJobNotReady(ctx, self.UserCred) self.OnDeployGuestFail(ctx, guest, fmt.Errorf("deploy backup failed %s", data)) } @@ -88,10 +95,12 @@ func (self *GuestDeployBackupTask) OnInit(ctx context.Context, obj db.IStandalon } func (self *GuestDeployBackupTask) OnDeployGuestComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + guest.SetGuestBackupMirrorJobNotReady(ctx, self.UserCred) self.SetStageComplete(ctx, nil) } func (self *GuestDeployBackupTask) OnDeployGuestCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + guest.SetGuestBackupMirrorJobNotReady(ctx, self.UserCred) guest.SetStatus(self.UserCred, api.VM_DEPLOYING_BACKUP_FAILED, data.String()) 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 2142224460..a5f8bc25db 100644 --- a/pkg/compute/tasks/ha_guest_start_task.go +++ b/pkg/compute/tasks/ha_guest_start_task.go @@ -38,10 +38,40 @@ func (self *HAGuestStartTask) OnInit( ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject, ) { guest := obj.(*models.SGuest) + host := models.HostManager.FetchHostById(guest.BackupHostId) + if host.HostStatus != api.HOST_ONLINE { + // request start master guest + self.GuestStartTask.OnInit(ctx, guest, nil) + } else { + self.RequestStopBackupGuest(ctx, guest) + } +} + +func (self *HAGuestStartTask) RequestStopBackupGuest(ctx context.Context, guest *models.SGuest) { + host := models.HostManager.FetchHostById(guest.BackupHostId) + self.SetStage("OnBackupGuestStopComplete", nil) + guest.SetStatus(self.UserCred, api.VM_BACKUP_STOPING, "HAGuestStartTask") + err := guest.GetDriver().RequestStopOnHost(ctx, guest, host, self, false) + if err != nil { + guest.SetStatus(self.UserCred, api.VM_BACKUP_START_FAILED, err.Error()) + self.SetStageFailed(ctx, nil) + } +} + +func (self *HAGuestStartTask) OnBackupGuestStopComplete( + ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, +) { db.OpsLog.LogEvent(guest, db.ACT_STARTING, nil, self.UserCred) self.RequestStartBacking(ctx, guest) } +func (self *HAGuestStartTask) OnBackupGuestStopCompleteFailed( + ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, +) { + guest.SetStatus(self.UserCred, api.VM_BACKUP_START_FAILED, data.String()) + self.SetStageFailed(ctx, data) +} + func (self *HAGuestStartTask) RequestStartBacking(ctx context.Context, guest *models.SGuest) { self.SetStage("OnStartBackupGuestComplete", nil) host := models.HostManager.FetchHostById(guest.BackupHostId) @@ -57,7 +87,7 @@ func (self *HAGuestStartTask) RequestStartBacking(ctx context.Context, guest *mo func (self *HAGuestStartTask) OnStartBackupGuestComplete( ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, ) { - if data != nil { + if data != nil && !jsonutils.QueryBoolean(data, "is_running", false) { nbdServerPort, err := data.Int("nbd_server_port") if err == nil { backupHost := models.HostManager.FetchHostById(guest.BackupHostId) @@ -69,12 +99,17 @@ func (self *HAGuestStartTask) OnStartBackupGuestComplete( return } } - + if err := guest.ResetGuestQuorumChildIndex(ctx, self.UserCred); err != nil { + self.OnStartBackupGuestCompleteFailed(ctx, guest, jsonutils.NewString(fmt.Sprintf("failed reset quorum child index: %s", err))) + return + } + guest.SetBackupGuestStatus(self.UserCred, api.VM_RUNNING, "on start backup guest complete") self.RequestStart(ctx, guest) } func (self *HAGuestStartTask) OnStartBackupGuestCompleteFailed( ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, ) { + guest.SetBackupGuestStatus(self.UserCred, api.VM_START_FAILED, data.String()) self.OnStartCompleteFailed(ctx, guest, data) } diff --git a/pkg/compute/tasks/ha_guest_stop_task.go b/pkg/compute/tasks/ha_guest_stop_task.go index 2de216ebaf..9960bfb81c 100644 --- a/pkg/compute/tasks/ha_guest_stop_task.go +++ b/pkg/compute/tasks/ha_guest_stop_task.go @@ -20,6 +20,8 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/apis/compute" + api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/compute/models" ) @@ -36,6 +38,11 @@ func (self *HAGuestStopTask) OnGuestStopTaskComplete( ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, ) { host := models.HostManager.FetchHostById(guest.BackupHostId) + if host.HostStatus != api.HOST_ONLINE { + self.GuestStopTask.OnGuestStopTaskComplete(ctx, guest, data) + return + } + self.SetStage("OnSlaveGuestStopTaskComplete", nil) err := guest.GetDriver().RequestStopOnHost(ctx, guest, host, self, true) if err != nil { @@ -54,5 +61,6 @@ func (self *HAGuestStopTask) OnSlaveGuestStopTaskComplete( func (self *HAGuestStopTask) OnSlaveGuestStopTaskCompleteFailed( ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, ) { + guest.SetBackupGuestStatus(self.UserCred, compute.VM_STOP_FAILED, data.String()) self.OnGuestStopTaskCompleteFailed(ctx, guest, data) } diff --git a/pkg/hostman/guestman/guesthandlers/guesthandler.go b/pkg/hostman/guestman/guesthandlers/guesthandler.go index 980fcbddfe..905fe70951 100644 --- a/pkg/hostman/guestman/guesthandlers/guesthandler.go +++ b/pkg/hostman/guestman/guesthandlers/guesthandler.go @@ -62,39 +62,40 @@ func AddGuestTaskHandler(prefix string, app *appsrv.Application) { auth.Authenticate(deleteGuest)) for action, f := range map[string]actionFunc{ - "create": guestCreate, - "deploy": guestDeploy, - "rebuild": guestRebuild, - "start": guestStart, - "stop": guestStop, - "monitor": guestMonitor, - "sync": guestSync, - "suspend": guestSuspend, - "io-throttle": guestIoThrottle, - "snapshot": guestSnapshot, - "delete-snapshot": guestDeleteSnapshot, - "reload-disk-snapshot": guestReloadDiskSnapshot, - "src-prepare-migrate": guestSrcPrepareMigrate, - "dest-prepare-migrate": guestDestPrepareMigrate, - "live-migrate": guestLiveMigrate, - "resume": guestResume, - "drive-mirror": guestDriveMirror, - "hotplug-cpu-mem": guestHotplugCpuMem, - "cancel-block-jobs": guestCancelBlockJobs, - "create-from-libvirt": guestCreateFromLibvirt, - "create-form-esxi": guestCreateFromEsxi, - "open-forward": guestOpenForward, - "list-forward": guestListForward, - "close-forward": guestCloseForward, - "storage-clone-disk": guestStorageCloneDisk, - "live-change-disk": guestLiveChangeDisk, - "cpuset": guestCPUSet, - "cpuset-remove": guestCPUSetRemove, - "memory-snapshot": guestMemorySnapshot, - "memory-snapshot-reset": guestMemorySnapshotReset, - "qga-set-password": qgaGuestSetPassword, - "qga-guest-ping": qgaGuestPing, - "qga-command": qgaCommand, + "create": guestCreate, + "deploy": guestDeploy, + "rebuild": guestRebuild, + "start": guestStart, + "stop": guestStop, + "monitor": guestMonitor, + "sync": guestSync, + "suspend": guestSuspend, + "io-throttle": guestIoThrottle, + "snapshot": guestSnapshot, + "delete-snapshot": guestDeleteSnapshot, + "reload-disk-snapshot": guestReloadDiskSnapshot, + "src-prepare-migrate": guestSrcPrepareMigrate, + "dest-prepare-migrate": guestDestPrepareMigrate, + "live-migrate": guestLiveMigrate, + "resume": guestResume, + "block-replication": guestBlockReplication, + "hotplug-cpu-mem": guestHotplugCpuMem, + "cancel-block-jobs": guestCancelBlockJobs, + "cancel-block-replication": guestCancelBlockReplication, + "create-from-libvirt": guestCreateFromLibvirt, + "create-form-esxi": guestCreateFromEsxi, + "open-forward": guestOpenForward, + "list-forward": guestListForward, + "close-forward": guestCloseForward, + "storage-clone-disk": guestStorageCloneDisk, + "live-change-disk": guestLiveChangeDisk, + "cpuset": guestCPUSet, + "cpuset-remove": guestCPUSetRemove, + "memory-snapshot": guestMemorySnapshot, + "memory-snapshot-reset": guestMemorySnapshotReset, + "qga-set-password": qgaGuestSetPassword, + "qga-guest-ping": qgaGuestPing, + "qga-command": qgaCommand, } { app.AddHandler("POST", fmt.Sprintf("%s/%s//%s", prefix, keyWord, action), @@ -474,7 +475,7 @@ func guestResume(ctx context.Context, userCred mcclient.TokenCredential, sid str // return nil, nil // } -func guestDriveMirror(ctx context.Context, userCred mcclient.TokenCredential, sid string, body jsonutils.JSONObject) (interface{}, error) { +func guestBlockReplication(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) } @@ -487,8 +488,7 @@ func guestDriveMirror(ctx context.Context, userCred mcclient.TokenCredential, si if err != nil { return nil, httperrors.NewInputParameterError("failed unmarshal desc %s", err) } - - hostutils.DelayTaskWithoutReqctx(ctx, guestman.GetGuestManager().StartDriveMirror, + hostutils.DelayTaskWithoutReqctx(ctx, guestman.GetGuestManager().StartBlockReplication, &guestman.SDriverMirror{ Sid: sid, NbdServerUri: backupNbdServerUri, @@ -505,6 +505,16 @@ func guestCancelBlockJobs(ctx context.Context, userCred mcclient.TokenCredential return nil, nil } +func guestCancelBlockReplication( + 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) + } + hostutils.DelayTaskWithoutReqctx(ctx, guestman.GetGuestManager().CancelBlockReplication, sid) + return nil, nil +} + func guestHotplugCpuMem(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 21df39ce07..9cc9d202d8 100644 --- a/pkg/hostman/guestman/guestman.go +++ b/pkg/hostman/guestman/guestman.go @@ -742,20 +742,23 @@ func (m *SGuestManager) StatusWithBlockJobsCount(ctx context.Context, params int guest, _ := m.GetServer(sid) 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("Block job missing") }) + 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() - body := jsonutils.NewDict() + body.Set("status", jsonutils.NewString(status)) - body.Set("block_jobs_count", jsonutils.NewInt(int64(blockJobsCount))) hostutils.TaskComplete(ctx, body) } if guest.Monitor == nil && !guest.IsStopping() { @@ -1220,17 +1223,29 @@ func (m *SGuestManager) OnlineResizeDisk(ctx context.Context, sid string, diskId // } -func (m *SGuestManager) StartDriveMirror(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { +func (m *SGuestManager) StartBlockReplication(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { mirrorParams, ok := params.(*SDriverMirror) if !ok { return nil, hostutils.ParamsError } + + nbdOpts := strings.Split(mirrorParams.NbdServerUri, ":") + if len(nbdOpts) != 3 { + return nil, fmt.Errorf("Nbd url is not vaild %s", mirrorParams.NbdServerUri) + } guest, _ := m.GetServer(mirrorParams.Sid) // TODO: check desc if err := guest.SaveSourceDesc(mirrorParams.Desc); err != nil { return nil, err } - task := NewDriveMirrorTask(ctx, guest, mirrorParams.NbdServerUri, "top", true, nil) + onSucc := func() { + if err := guest.updateChildIndex(); err != nil { + hostutils.TaskFailed(ctx, err.Error()) + return + } + hostutils.TaskComplete(ctx, nil) + } + task := NewGuestBlockReplicationTask(ctx, guest, nbdOpts[1], nbdOpts[2], "top", onSucc, nil) task.Start() return nil, nil } @@ -1256,6 +1271,27 @@ func (m *SGuestManager) CancelBlockJobs(ctx context.Context, params interface{}) return nil, nil } +func (m *SGuestManager) CancelBlockReplication(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { + sid, ok := params.(string) + if !ok { + return nil, hostutils.ParamsError + } + status := m.getStatus(sid) + if status == GUSET_STOPPED { + hostutils.TaskComplete(ctx, nil) + return nil, nil + } + defer func() { + if r := recover(); r != nil { + log.Errorf("STACK: %v \n %s", r, debug.Stack()) + hostutils.TaskFailed(ctx, fmt.Sprintf("recover: %v", r)) + } + }() + guest, _ := m.GetServer(sid) + NewCancelBlockReplicationTask(ctx, guest).Start() + return nil, nil +} + func (m *SGuestManager) HotplugCpuMem(ctx context.Context, params interface{}) (jsonutils.JSONObject, error) { hotplugParams, ok := params.(*SGuestHotplugCpuMem) if !ok { diff --git a/pkg/hostman/guestman/guesttasks.go b/pkg/hostman/guestman/guesttasks.go index d118f13eb4..588b0c8460 100644 --- a/pkg/hostman/guestman/guesttasks.go +++ b/pkg/hostman/guestman/guesttasks.go @@ -503,7 +503,7 @@ func (d *SGuestDiskSyncTask) startAddDisk(disk *desc.SGuestDisk) { bus = fmt.Sprintf("ide.%d", diskIndex) } // drive_add bus is a placeholder - d.guest.Monitor.DriveAdd(bus, params, func(result string) { d.onAddDiskSucc(disk, result, pciRoot) }) + d.guest.Monitor.DriveAdd(bus, "", params, func(result string) { d.onAddDiskSucc(disk, result, pciRoot) }) } func (d *SGuestDiskSyncTask) onAddDiskSucc(disk *desc.SGuestDisk, results string, pciRoot *desc.PCIController) { @@ -1948,7 +1948,7 @@ func (s *SDriveMirrorTask) startMirror(res string) { log.Infof("mirror block replication supported") } if s.index < len(s.Desc.Disks) { - target := fmt.Sprintf("%s:exportname=drive_%d", s.nbdUri, s.index) + target := fmt.Sprintf("%s:exportname=drive_%d_backend", s.nbdUri, s.index) s.Monitor.DriveMirror(s.startMirror, fmt.Sprintf("drive_%d", s.index), target, s.syncMode, "", true, blockReplication) s.index += 1 @@ -1961,6 +1961,85 @@ func (s *SDriveMirrorTask) startMirror(res string) { } } +type SGuestBlockReplicationTask struct { + *SKVMGuestInstance + + ctx context.Context + nbdHost string + nbdPort string + onSucc func() + onFail func(string) + syncMode string + index int +} + +func NewGuestBlockReplicationTask( + ctx context.Context, s *SKVMGuestInstance, + nbdHost, nbdPort, syncMode string, onSucc func(), onFail func(string), +) *SGuestBlockReplicationTask { + return &SGuestBlockReplicationTask{ + SKVMGuestInstance: s, + ctx: ctx, + nbdHost: nbdHost, + nbdPort: nbdPort, + syncMode: syncMode, + onSucc: onSucc, + onFail: onFail, + } +} + +func (s *SGuestBlockReplicationTask) Start() { + s.onXBlockdevChange("") +} + +func (s *SGuestBlockReplicationTask) onXBlockdevChange(res string) { + if len(res) > 0 { + log.Errorf("SGuestBlockReplicationTask onXBlockdevChange %s", res) + if s.onFail != nil { + s.onFail(res) + } else { + hostutils.TaskFailed(s.ctx, res) + } + return + } + + disks := s.Desc.Disks + if s.index < len(disks) { + diskIndex := disks[s.index].Index + drive := fmt.Sprintf("drive_%d", diskIndex) + node := fmt.Sprintf("node_%d", diskIndex) + + s.Monitor.DriveAdd("", "buddy", map[string]string{ + "file.driver": "nbd", "file.host": s.nbdHost, "file.port": s.nbdPort, + "file.export": drive, "node-name": node, + }, s.onNbdDriveAddSucc(drive, node)) + s.index += 1 + } else { + s.startDriveMirror() + } +} + +func (s *SGuestBlockReplicationTask) onNbdDriveAddSucc(parent, node string) monitor.StringCallback { + return func(res string) { + if len(res) > 0 { + log.Errorf("SGuestBlockReplicationTask onNbdDriveAddSucc %s", res) + if s.onFail != nil { + s.onFail(res) + } else { + hostutils.TaskFailed(s.ctx, res) + } + return + } + + s.Monitor.XBlockdevChange(parent, node, "", s.onXBlockdevChange) + } +} + +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 **/ @@ -2289,6 +2368,38 @@ func (task *SGuestBlockIoThrottleTask) doIoThrottle(drivers []string) { } } +type CancelBlockReplication struct { + SCancelBlockJobs +} + +func NewCancelBlockReplicationTask(ctx context.Context, guest *SKVMGuestInstance) *CancelBlockReplication { + return &CancelBlockReplication{SCancelBlockJobs{guest, ctx}} +} + +func (task *CancelBlockReplication) Start() { + // start remove child node of block device + disks := task.Desc.Disks + for i := 0; i < len(disks); i++ { + diskIndex := disks[i].Index + drive := fmt.Sprintf("drive_%d", diskIndex) + node := fmt.Sprintf("node_%d", diskIndex) + child := fmt.Sprintf("children.%d", task.getQuorumChildIndex()) + task.Monitor.XBlockdevChange(drive, "", child, func(res string) { + if len(res) > 0 { + log.Errorf("failed remove child %s for parent %s: %s", drive, node, res) + return + } + task.Monitor.DriveDel(node, func(res string) { + if len(res) > 0 { + log.Errorf("failed remove drive %s: %s", node, res) + return + } + }) + }) + } + task.SCancelBlockJobs.Start() +} + type SCancelBlockJobs struct { *SKVMGuestInstance @@ -2327,7 +2438,7 @@ func (task *SCancelBlockJobs) StartCancelBlockJobs(drivers []string) { } task.StartCancelBlockJobs(drivers) } - task.Monitor.CancelBlockJob(driver, true, onCancelBlockJob) + task.Monitor.CancelBlockJob(driver, false, onCancelBlockJob) } else { task.taskComplete() } diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index f8a62e48b6..2bdf5ba435 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -24,6 +24,7 @@ import ( "regexp" "strconv" "strings" + "sync/atomic" "syscall" "time" @@ -94,6 +95,7 @@ type SKVMInstanceRuntime struct { stopping bool needSyncStreamDisks bool blockJobTigger map[string]chan struct{} + quorumFailed int32 StartupTask *SGuestResumeTask MigrateTask *SGuestLiveMigrateTask @@ -486,6 +488,7 @@ 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 { @@ -613,9 +616,7 @@ func (s *SKVMGuestInstance) ImportServer(pendingDelete bool) { action = "suspend" } log.Infof("%s is %s, pending_delete=%t", s.GetName(), action, pendingDelete) - if !s.IsSlave() { - s.SyncStatus("") - } + s.SyncStatus("") } } @@ -745,53 +746,58 @@ func (s *SKVMGuestInstance) StartMonitor(ctx context.Context, cb func()) error { } func (s *SKVMGuestInstance) onReceiveQMPEvent(event *monitor.Event) { - switch { - case event.Event == `"BLOCK_JOB_READY"`: + switch event.Event { + case `"BLOCK_JOB_READY"`, `"BLOCK_JOB_COMPLETED"`: s.eventBlockJobReady(event) - case event.Event == `"BLOCK_JOB_ERROR"`: - s.SyncMirrorJobFailed("BLOCK_JOB_ERROR") - case event.Event == `"BLOCK_JOB_COMPLETED"`: - s.eventBlockJobCompleted(event) - case event.Event == `"GUEST_PANICKED"`: + case `"BLOCK_JOB_ERROR"`: + s.eventBlockJobError(event) + case `"GUEST_PANICKED"`: s.eventGuestPaniced(event) - case event.Event == `"STOP"`: - if s.MigrateTask != nil { - s.MigrateTask.onMigrateReceivedStopEvent() - } + case `"STOP"`: + s.eventGuestStop() + case `"QUORUM_REPORT_BAD"`: + s.eventQuorumReportBad(event) } } -func (s *SKVMGuestInstance) eventBlockJobCompleted(event *monitor.Event) { - itype, ok := event.Data["type"] - if !ok { - log.Errorf("BLOCK_JOB_COMPLETED missing event type") - return +func (s *SKVMGuestInstance) eventBlockJobError(event *monitor.Event) { + s.SyncMirrorJobFailed(event.String()) +} + +func (s *SKVMGuestInstance) eventGuestStop() { + if s.MigrateTask != nil { + // migrating complete + s.MigrateTask.onMigrateReceivedStopEvent() } - // only dealwith event type mirror - stype, _ := itype.(string) - if stype != "mirror" { + hostutils.UpdateServerProgress(context.Background(), s.Id, 0.0, 0) +} + +func (s *SKVMGuestInstance) eventQuorumReportBad(event *monitor.Event) { + if !atomic.CompareAndSwapInt32(&s.quorumFailed, 0, 1) { return } - 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 complete 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 - } - diskId := disks[diskIndex].DiskId - if c, ok := s.blockJobTigger[diskId]; ok { - c <- struct{}{} + for i := 0; i < len(disks); i++ { + diskIndex := disks[i].Index + drive := fmt.Sprintf("drive_%d", diskIndex) + node := fmt.Sprintf("node_%d", diskIndex) + child := fmt.Sprintf("children.%d", s.getQuorumChildIndex()) + s.Monitor.XBlockdevChange(drive, "", child, func(res string) { + if len(res) > 0 { + log.Errorf("On QUORUM_REPORT_BAD failed remove child %s for parent %s: %s", drive, node, res) + return + } + s.Monitor.DriveDel(node, func(res string) { + if len(res) > 0 { + log.Errorf("On QUORUM_REPORT_BAD failed remove drive %s: %s", node, res) + return + } + }) + }) } + + s.SyncMirrorJobFailed(event.String()) } func (s *SKVMGuestInstance) eventGuestPaniced(event *monitor.Event) { @@ -816,7 +822,7 @@ func (s *SKVMGuestInstance) eventGuestPaniced(event *monitor.Event) { func (s *SKVMGuestInstance) eventBlockJobReady(event *monitor.Event) { itype, ok := event.Data["type"] if !ok { - log.Errorf("BLOCK_JOB_READY missing event type") + log.Errorf("block job missing event type") return } // only dealwith event type mirror @@ -828,12 +834,23 @@ func (s *SKVMGuestInstance) eventBlockJobReady(event *monitor.Event) { if s.IsMaster() { // has backup server mirrorStatus := s.MirrorJobStatus() if mirrorStatus.IsSucc() { - _, err := hostutils.UpdateServerStatus(context.Background(), s.GetId(), api.VM_RUNNING, s.GetPowerStates(), "BLOCK_JOB_READY") - if err != nil { - log.Errorf("onReceiveQMPEvent update server status error: %s", err) + 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) + if err != nil { + log.Errorf("onReceiveQMPEvent update server status error: %s", err) + time.Sleep(3 * time.Second) + } else { + break + } } } else if mirrorStatus.IsFailed() { - s.SyncMirrorJobFailed("Block job missing") + s.SyncMirrorJobFailed("drive-mirror job failed") } } else { iDevice, ok := event.Data["device"] @@ -986,7 +1003,7 @@ func (s *SKVMGuestInstance) getMemoryDevices(ctx context.Context, pciInfoList [] } func (s *SKVMGuestInstance) guestRun(ctx context.Context) { - if s.LiveMigrateDestPort != nil && ctx != nil { + if s.LiveMigrateDestPort != nil && ctx != nil && !s.IsSlave() { // dest migrate guest body := jsonutils.NewDict() body.Set("live_migrate_dest_port", jsonutils.NewInt(*s.LiveMigrateDestPort)) @@ -1005,14 +1022,6 @@ func (s *SKVMGuestInstance) guestRun(ctx context.Context) { } else { if s.IsMaster() { s.startDiskBackupMirror(ctx) - if ctx != nil && len(appctx.AppContextTaskId(ctx)) > 0 { - s.DoResumeTask(ctx, false) - } else { - if options.HostOptions.SetVncPassword { - s.SetVncPassword() - } - s.OnResumeSyncMetadataInfo() - } } else { s.DoResumeTask(ctx, true) } @@ -1026,9 +1035,7 @@ func (s *SKVMGuestInstance) onMonitorDisConnect(err error) { log.Errorf("Guest %s on Monitor Disconnect reason: %v", s.Id, err) s.CleanStartupTask() s.scriptStop() - if !s.IsSlave() && s.LiveMigrateDestPort == nil { - s.SyncStatus(fmt.Sprintf("monitor disconnect %v", err)) - } + s.SyncStatus(fmt.Sprintf("monitor disconnect %v", err)) if s.guestAgent != nil { s.guestAgent.Close() s.guestAgent = nil @@ -1039,26 +1046,31 @@ func (s *SKVMGuestInstance) onMonitorDisConnect(err error) { func (s *SKVMGuestInstance) startDiskBackupMirror(ctx context.Context) { if ctx == nil || len(appctx.AppContextTaskId(ctx)) == 0 { - status := api.VM_RUNNING - mirrorStatus := s.MirrorJobStatus() - if mirrorStatus.InProcess() { - status = api.VM_BLOCK_STREAM - } else if mirrorStatus.IsFailed() { - status = api.VM_BLOCK_STREAM_FAIL - s.SyncMirrorJobFailed("mirror job missing") - } - hostutils.UpdateServerStatus(context.Background(), s.GetId(), status, s.GetPowerStates(), "") + s.DoResumeTask(ctx, true) } else { nbdUri, ok := s.Desc.Metadata["backup_nbd_server_uri"] if !ok { hostutils.TaskFailed(ctx, "Missing dest nbd location") + return } - + nbdOpts := strings.Split(nbdUri, ":") + if len(nbdOpts) != 3 { + hostutils.TaskFailed(ctx, fmt.Sprintf("Nbd uri is not vaild %s", nbdUri)) + return + } + s.quorumFailed = 0 onSucc := func() { - cb := func(res string) { log.Infof("On backup mirror server(%s) resume start", s.Id) } - s.Monitor.SimpleCommand("cont", cb) + if err := s.updateChildIndex(); err != nil { + hostutils.TaskFailed(ctx, err.Error()) + return + } + s.DoResumeTask(ctx, true) } - NewDriveMirrorTask(ctx, s, nbdUri, "top", true, onSucc).Start() + onFail := func(res string) { + s.SyncMirrorJobFailed(res) + s.DoResumeTask(ctx, true) + } + NewGuestBlockReplicationTask(ctx, s, nbdOpts[1], nbdOpts[2], "top", onSucc, onFail).Start() } } @@ -1107,18 +1119,25 @@ func (s *SKVMGuestInstance) DiskCount() int { return len(s.Desc.Disks) } -type MirrorJob int +type MirrorJob struct { + mirrorJobStatus int + blockJobsCount int +} func (ms MirrorJob) IsSucc() bool { - return ms == 1 + return ms.mirrorJobStatus == 1 } func (ms MirrorJob) IsFailed() bool { - return ms == -1 + return ms.mirrorJobStatus == -1 } func (ms MirrorJob) InProcess() bool { - return ms == 0 + return ms.mirrorJobStatus == 0 +} + +func (ms MirrorJob) BlockJobsCount() int { + return ms.blockJobsCount } func (s *SKVMGuestInstance) MirrorJobStatus() MirrorJob { @@ -1128,21 +1147,27 @@ func (s *SKVMGuestInstance) MirrorJobStatus() MirrorJob { }) select { case <-time.After(time.Second * 3): - return 0 + return MirrorJob{0, -1} case v := <-res: - if len(v) >= s.DiskCount() { - mirrorSuccCount := 0 - for _, job := range v { - if job.Type == "mirror" && job.Status == "ready" { - mirrorSuccCount += 1 - } + mirrorJobCount := 0 + failedJobCount := 0 + for _, job := range v { + if job.Type != "mirror" { + continue } - if mirrorSuccCount == s.DiskCount() { - return 1 + if job.IoStatus != "ok" { + failedJobCount += 1 } - return 0 + mirrorJobCount += 1 + } + if failedJobCount > 0 { + return MirrorJob{-1, len(v)} + } + if mirrorJobCount == 0 { + return MirrorJob{1, len(v)} + } else { + return MirrorJob{0, len(v)} } - return -1 } } @@ -1237,12 +1262,18 @@ func (s *SKVMGuestInstance) SyncStatus(reason string) { s.Monitor.GetBlockJobCounts(s.CheckBlockOrRunning) return } - var status = "ready" + var status = api.VM_READY if s.IsSuspend() { - status = "suspend" + status = api.VM_SUSPEND + } + statusInput := &apis.PerformStatusInput{ + Status: status, + Reason: reason, + IsSlave: s.IsSlave(), + PowerStates: s.GetPowerStates(), } - if _, err := hostutils.UpdateServerStatus(context.Background(), s.Id, status, s.GetPowerStates(), reason); err != nil { + if _, err := hostutils.UpdateServerStatus(context.Background(), s.Id, statusInput); err != nil { log.Errorf("failed update guest status %s", err) } } @@ -1257,6 +1288,7 @@ func (s *SKVMGuestInstance) GetPowerStates() string { func (s *SKVMGuestInstance) CheckBlockOrRunning(jobs int) { var status = api.VM_RUNNING + if jobs > 0 { if s.IsMaster() { mirrorStatus := s.MirrorJobStatus() @@ -1264,14 +1296,19 @@ func (s *SKVMGuestInstance) CheckBlockOrRunning(jobs int) { status = api.VM_BLOCK_STREAM } else if mirrorStatus.IsFailed() { status = api.VM_BLOCK_STREAM_FAIL - s.SyncMirrorJobFailed("Block job missing") + s.SyncMirrorJobFailed("drive-mirror job failed") } } else { // TODO: check block jobs ready status = api.VM_BLOCK_STREAM } } - _, err := hostutils.UpdateServerStatus(context.Background(), s.Id, status, s.GetPowerStates(), "") + var statusInput = &apis.PerformStatusInput{ + Status: status, + BlockJobsCount: jobs, + PowerStates: s.GetPowerStates(), + } + _, err := hostutils.UpdateServerStatus(context.Background(), s.Id, statusInput) if err != nil { log.Errorln(err) } @@ -2144,6 +2181,15 @@ func (s *SKVMGuestInstance) SyncMetadata(meta *jsonutils.JSONDict) error { return nil } +func (s *SKVMGuestInstance) updateChildIndex() error { + idx := s.getQuorumChildIndex() + 1 + s.Desc.Metadata[api.QUORUM_CHILD_INDEX] = strconv.Itoa(int(idx)) + s.SaveLiveDesc(s.Desc) + meta := jsonutils.NewDict() + meta.Set(api.QUORUM_CHILD_INDEX, jsonutils.NewInt(idx)) + return s.SyncMetadata(meta) +} + func (s *SKVMGuestInstance) SetVncPassword() { password := seclib.RandomPassword(8) s.VncPassword = password diff --git a/pkg/hostman/guestman/qemu-kvmhelper.go b/pkg/hostman/guestman/qemu-kvmhelper.go index bfd37b586b..a27612a334 100644 --- a/pkg/hostman/guestman/qemu-kvmhelper.go +++ b/pkg/hostman/guestman/qemu-kvmhelper.go @@ -38,7 +38,10 @@ import ( qemucerts "yunion.io/x/onecloud/pkg/hostman/guestman/qemu/certs" "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" ) @@ -253,6 +256,14 @@ func (s *SKVMGuestInstance) disablePvpanicDev() bool { return s.Desc.Metadata["disable_pvpanic"] == "true" } +func (s *SKVMGuestInstance) getQuorumChildIndex() int64 { + if sidx, ok := s.Desc.Metadata[api.QUORUM_CHILD_INDEX]; ok { + idx, _ := strconv.ParseInt(sidx, 10, 0) + return idx + } + return 0 +} + func (s *SKVMGuestInstance) getNicUpScriptPath(nic *desc.SGuestNetwork) string { dev := s.manager.GetHost().GetBridgeDev(nic.Bridge) return path.Join(s.HomeDir(), fmt.Sprintf("if-up-%s-%s.sh", dev.Bridge(), nic.Ifname)) @@ -498,6 +509,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 { + return "", err + } + } + qemuOpts, err := qemu.GenerateStartOptions(input) if err != nil { return "", errors.Wrap(err, "GenerateStartCommand") @@ -516,6 +533,49 @@ echo $CMD` return cmd, nil } +func (s *SKVMGuestInstance) slaveDiskPrepare(input *qemu.GenerateStartOptionsInput) 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) + } + 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") + } + } + } + return nil +} + func (s *SKVMGuestInstance) parseCmdline(input string) (*qemutils.Cmdline, []qemutils.Option, error) { cl, err := qemutils.NewCmdline(input) if err != nil { diff --git a/pkg/hostman/guestman/qemu/generate.go b/pkg/hostman/guestman/qemu/generate.go index cb51694660..0d074a3d89 100644 --- a/pkg/hostman/guestman/qemu/generate.go +++ b/pkg/hostman/guestman/qemu/generate.go @@ -202,26 +202,57 @@ func getMonitorOptions(drvOpt QemuOptions, input *Monitor) []string { return opts } -func generateDisksOptions(drvOpt QemuOptions, disks []*desc.SGuestDisk, isEncrypt bool) []string { +func generateDisksOptions(drvOpt QemuOptions, disks []*desc.SGuestDisk, isEncrypt, isSlave, isMaster bool) []string { opts := make([]string, 0) for _, disk := range disks { - opts = append(opts, - getDiskDriveOption(drvOpt, disk, isEncrypt), - getDiskDeviceOption(drvOpt, disk), - ) + 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, getDiskDeviceOption(drvOpt, disk)) } return opts } -func getDiskDriveOption(drvOpt QemuOptions, disk *desc.SGuestDisk, isEncrypt bool) string { +func getMasterDiskDriveOption(drvOpt QemuOptions, disk *desc.SGuestDisk, isEncrypt bool) string { + format := disk.Format + diskIndex := disk.Index + cacheMode := disk.CacheMode + aioMode := disk.AioMode + opt := "if=none,driver=quorum,read-pattern=fifo,is-backup-mode=on,vote-threshold=1" + opt += fmt.Sprintf(",id=drive_%d", diskIndex) + opt += fmt.Sprintf(",cache=%s", cacheMode) + if isLocalStorage(disk) { + opt += fmt.Sprintf(",aio=%s", aioMode) + } + opt += fmt.Sprintf(",children.0.file.filename=$DISK_%d", diskIndex) + if format == "raw" { + opt += ",children.0.file.format=raw" + } + return drvOpt.Drive(opt) +} + +func getDiskDriveOption(drvOpt QemuOptions, disk *desc.SGuestDisk, isEncrypt, isSlave bool) string { format := disk.Format diskIndex := disk.Index cacheMode := disk.CacheMode aioMode := disk.AioMode - opt := fmt.Sprintf("file=$DISK_%d", diskIndex) - opt += ",if=none" - opt += fmt.Sprintf(",id=drive_%d", diskIndex) + 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) + } + if len(format) == 0 || format == "qcow2" { // pass # qemu will automatically detect image format } else if format == "raw" { @@ -686,7 +717,8 @@ func GenerateStartOptions( opts = append(opts, generatePCIDeviceOption(input.GuestDesc.PvScsi.PCIDevice)) } // generate disk options - opts = append(opts, generateDisksOptions(drvOpt, input.GuestDesc.Disks, isEncrypt)...) + opts = append(opts, generateDisksOptions( + drvOpt, input.GuestDesc.Disks, isEncrypt, input.GuestDesc.IsSlave, input.GuestDesc.IsMaster)...) // cdrom opts = append(opts, generateCdromOptions(drvOpt, input.GuestDesc.Cdroms)...) diff --git a/pkg/hostman/hostutils/hostutils.go b/pkg/hostman/hostutils/hostutils.go index 9c37f9dc98..34d987081a 100644 --- a/pkg/hostman/hostutils/hostutils.go +++ b/pkg/hostman/hostutils/hostutils.go @@ -22,6 +22,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/apis" "yunion.io/x/onecloud/pkg/appctx" "yunion.io/x/onecloud/pkg/appsrv" "yunion.io/x/onecloud/pkg/cloudcommon/consts" @@ -142,14 +143,8 @@ func RemoteStoragecacheCacheImage(ctx context.Context, storagecacheId, imageId, storagecacheId, imageId, query, params) } -func UpdateServerStatus(ctx context.Context, sid, status, powerStates, reason string) (jsonutils.JSONObject, error) { - var stats = jsonutils.NewDict() - stats.Set("status", jsonutils.NewString(status)) - stats.Set("power_states", jsonutils.NewString(powerStates)) - if len(reason) > 0 { - stats.Set("reason", jsonutils.NewString(reason)) - } - return modules.Servers.PerformAction(GetComputeSession(ctx), sid, "status", stats) +func UpdateServerStatus(ctx context.Context, sid string, statusInput *apis.PerformStatusInput) (jsonutils.JSONObject, error) { + return modules.Servers.PerformAction(GetComputeSession(ctx), sid, "status", jsonutils.Marshal(statusInput)) } func UpdateServerProgress(ctx context.Context, sid string, progress, progressMbps float64) (jsonutils.JSONObject, error) { diff --git a/pkg/hostman/monitor/hmp.go b/pkg/hostman/monitor/hmp.go index 36834b91d9..40f31c5979 100644 --- a/pkg/hostman/monitor/hmp.go +++ b/pkg/hostman/monitor/hmp.go @@ -277,12 +277,20 @@ func (m *HmpMonitor) ObjectDel(idstr string, callback StringCallback) { m.Query(fmt.Sprintf("object_del %s", idstr), callback) } -func (m *HmpMonitor) DriveAdd(bus string, params map[string]string, callback StringCallback) { +func (m *HmpMonitor) XBlockdevChange(parent, node, child string, callback StringCallback) { + go callback("hmp not support command x-blockdev-change") +} + +func (m *HmpMonitor) DriveAdd(bus, node string, params map[string]string, callback StringCallback) { var paramsKvs = []string{} for k, v := range params { paramsKvs = append(paramsKvs, fmt.Sprintf("%s=%s", k, v)) } - m.Query(fmt.Sprintf("drive_add %s %s", bus, strings.Join(paramsKvs, ",")), callback) + cmd := "drive_add" + if len(node) > 0 { + cmd = fmt.Sprintf("drive_add -n %s", node) + } + m.Query(fmt.Sprintf("%s %s %s", cmd, bus, strings.Join(paramsKvs, ",")), callback) } func (m *HmpMonitor) DeviceAdd(dev string, params map[string]string, callback StringCallback) { diff --git a/pkg/hostman/monitor/monitor.go b/pkg/hostman/monitor/monitor.go index 1c5bf77d2c..cdc63f9a29 100644 --- a/pkg/hostman/monitor/monitor.go +++ b/pkg/hostman/monitor/monitor.go @@ -215,9 +215,10 @@ type Monitor interface { ObjectDel(idstr string, callback StringCallback) ObjectAdd(objectType string, params map[string]string, callback StringCallback) - DriveAdd(bus string, params map[string]string, callback StringCallback) + DriveAdd(bus, node string, params map[string]string, callback StringCallback) DeviceAdd(dev string, params map[string]string, callback StringCallback) + 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) BlockJobComplete(drive string, cb StringCallback) diff --git a/pkg/hostman/monitor/qmp.go b/pkg/hostman/monitor/qmp.go index ddfd9f7565..d9572ad65a 100644 --- a/pkg/hostman/monitor/qmp.go +++ b/pkg/hostman/monitor/qmp.go @@ -554,12 +554,37 @@ func (m *QmpMonitor) ObjectDel(idstr string, callback StringCallback) { m.HumanMonitorCommand(fmt.Sprintf("object_del %s", idstr), callback) } -func (m *QmpMonitor) DriveAdd(bus string, params map[string]string, callback StringCallback) { +func (m *QmpMonitor) XBlockdevChange(parent, node, child string, callback StringCallback) { + cb := func(res *Response) { + callback(m.actionResult(res)) + } + cmd := &Command{ + Execute: "x-blockdev-change", + } + args := map[string]interface{}{ + "parent": parent, + } + if len(node) > 0 { + args["node"] = node + } + if len(child) > 0 { + args["child"] = child + } + cmd.Args = args + m.Query(cmd, cb) +} + +func (m *QmpMonitor) DriveAdd(bus, node string, params map[string]string, callback StringCallback) { var paramsKvs = []string{} for k, v := range params { paramsKvs = append(paramsKvs, fmt.Sprintf("%s=%s", k, v)) } - cmd := fmt.Sprintf("drive_add %s %s", bus, strings.Join(paramsKvs, ",")) + cmd := "drive_add" + if len(node) > 0 { + cmd = fmt.Sprintf("drive_add -n %s", node) + } + + cmd = fmt.Sprintf("%s %s %s", cmd, bus, strings.Join(paramsKvs, ",")) m.HumanMonitorCommand(cmd, callback) // XXX: 同下 // var ( diff --git a/pkg/hostman/storageman/storage_local.go b/pkg/hostman/storageman/storage_local.go index df35ea20da..93498cc2f5 100644 --- a/pkg/hostman/storageman/storage_local.go +++ b/pkg/hostman/storageman/storage_local.go @@ -312,8 +312,45 @@ func (s *SLocalStorage) Detach() error { return nil } +func (s *SLocalStorage) deleteBackendFile(diskpath string, skipRecycle bool) error { + backendPath := diskpath + ".backend" + if !fileutils2.Exists(backendPath) { + return nil + } + disk, err := qemuimg.NewQemuImage(diskpath) + if err != nil { + return errors.Wrapf(err, "qemuimg.NewQemuImage(%s)", diskpath) + } + if disk.BackFilePath != backendPath { + return nil + } + + destDir := s.getRecyclePath() + if options.HostOptions.RecycleDiskfile && (!skipRecycle || options.HostOptions.AlwaysRecycleDiskfile) { + if err := procutils.NewCommand("mkdir", "-p", destDir).Run(); err != nil { + log.Errorf("Fail to mkdir %s for recycle: %s", destDir, err) + return err + } + backendDestFile := fmt.Sprintf("%s.%d", path.Base(backendPath), time.Now().Unix()) + log.Infof("Move deleted disk file %s to recycle %s", backendPath, destDir) + return procutils.NewCommand("mv", "-f", backendPath, path.Join(destDir, backendDestFile)).Run() + } else { + log.Infof("Delete disk file %s immediately", backendPath) + if options.HostOptions.ZeroCleanDiskData { + // try to zero clean files in subdir + zeroclean.ZeroDir(backendPath) + } + return procutils.NewCommand("rm", "-rf", backendPath).Run() + } +} + func (s *SLocalStorage) DeleteDiskfile(diskpath string, skipRecycle bool) error { log.Infof("Start Delete %s", diskpath) + + if err := s.deleteBackendFile(diskpath, skipRecycle); err != nil { + return err + } + if options.HostOptions.RecycleDiskfile && (!skipRecycle || options.HostOptions.AlwaysRecycleDiskfile) { var ( destDir = s.getRecyclePath() @@ -323,6 +360,7 @@ func (s *SLocalStorage) DeleteDiskfile(diskpath string, skipRecycle bool) error log.Errorf("Fail to mkdir %s for recycle: %s", destDir, err) return err } + log.Infof("Move deleted disk file %s to recycle %s", diskpath, destDir) return procutils.NewCommand("mv", "-f", diskpath, path.Join(destDir, destFile)).Run() } else { diff --git a/pkg/mcclient/options/compute/servers.go b/pkg/mcclient/options/compute/servers.go index 0711904bf8..ea95c58ac1 100644 --- a/pkg/mcclient/options/compute/servers.go +++ b/pkg/mcclient/options/compute/servers.go @@ -156,9 +156,8 @@ func (o *ServerDeleteBackupOptions) Params() (jsonutils.JSONObject, error) { } type ServerSwitchToBackupOptions struct { - ID string `help:"ID of the server" json:"-"` - PurgeBackup *bool `help:"Purge Guest Backup" json:"purge_backup"` - DeleteBackup *bool `help:"Delete Guest Backup" json:"delete_backup"` + ID string `help:"ID of the server" json:"-"` + AutoStart bool `help:"Start guest after switch to backup" json:"auto_start"` } func (o *ServerSwitchToBackupOptions) GetId() string { @@ -173,6 +172,19 @@ func (o *ServerSwitchToBackupOptions) Description() string { return "Switch geust master to backup host" } +type ServerCreateBackupOptions struct { + ID string `help:"ID of the server" json:"-"` + AutoStart bool `help:"Start guest after create backup guest" json:"auto_start"` +} + +func (o *ServerCreateBackupOptions) GetId() string { + return o.ID +} + +func (o *ServerCreateBackupOptions) Params() (jsonutils.JSONObject, error) { + return options.StructToParams(o) +} + type ServerShowOptions struct { options.BaseShowOptions `id->help:"ID or name of the server"` } diff --git a/pkg/util/logclient/consts.go b/pkg/util/logclient/consts.go index bbbdeb8c69..bc56d74f65 100644 --- a/pkg/util/logclient/consts.go +++ b/pkg/util/logclient/consts.go @@ -178,7 +178,8 @@ const ( ACT_FLUSH_INSTANCE = "flush_instance" - ACT_UPDATE_STATUS = "update_status" + ACT_UPDATE_STATUS = "update_status" + ACT_UPDATE_BACKUP_GUEST_STATUS = "update_backup_guest_status" ACT_UPDATE_PASSWORD = "update_password"