Merge pull request #15566 from wanyaoqi/automated-cherry-pick-of-#15565-upstream-master

Automated cherry pick of #15565: Automated cherry pick of #15473: fix(region,host): backup guest
This commit is contained in:
Zexi Li
2022-12-22 10:17:39 +08:00
committed by GitHub
30 changed files with 786 additions and 473 deletions
+3 -1
View File
@@ -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"
+5 -3
View File
@@ -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"
+2
View File
@@ -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"`
+2 -1
View File
@@ -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"`
+6 -5
View File
@@ -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"
+7 -2
View File
@@ -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 {
+73 -42
View File
@@ -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) {
+43 -110
View File
@@ -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)
}
+4 -11
View File
@@ -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)
-5
View File
@@ -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)
+41 -119
View File
@@ -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)
@@ -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)
}
@@ -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 {
+10 -1
View File
@@ -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)
}
+37 -2
View File
@@ -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)
}
+8
View File
@@ -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)
}
@@ -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/<sid>/%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)
+42 -6
View File
@@ -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 {
+114 -3
View File
@@ -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()
}
+135 -89
View File
@@ -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
+60
View File
@@ -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 {
+42 -10
View File
@@ -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)...)
+3 -8
View File
@@ -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) {
+10 -2
View File
@@ -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) {
+2 -1
View File
@@ -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)
+27 -2
View File
@@ -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 (
+38
View File
@@ -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 {
+15 -3
View File
@@ -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"`
}
+2 -1
View File
@@ -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"