From c0e2019acd868a907cfa5b94e7bb04fdffb2bff1 Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Thu, 9 Jan 2020 18:31:31 +0800 Subject: [PATCH] backup guest fix --- pkg/apis/compute/guest_const.go | 2 +- pkg/compute/guestdrivers/virtualization.go | 25 ++- pkg/compute/models/disks.go | 7 +- pkg/compute/models/guest_actions.go | 46 +++++- pkg/compute/tasks/disk_create_task.go | 73 +-------- pkg/compute/tasks/guest_backup_tasks.go | 18 +++ pkg/compute/tasks/guest_deploy_task.go | 52 +----- pkg/compute/tasks/guest_start_task.go | 73 +-------- pkg/compute/tasks/guest_stop_task.go | 21 +-- pkg/compute/tasks/ha_disk_create_task.go | 100 ++++++++++++ pkg/compute/tasks/ha_guest_deploy_task.go | 74 +++++++++ pkg/compute/tasks/ha_guest_start_task.go | 67 ++++++++ pkg/compute/tasks/ha_guest_stop_task.go | 44 +++++ pkg/hostman/guestman/guesttasks.go | 3 - pkg/hostman/guestman/qemu-kvm.go | 40 +++-- pkg/scheduler/api/sched.go | 28 +++- pkg/scheduler/core/context.go | 17 +- pkg/scheduler/core/generic_scheduler.go | 31 ++-- pkg/scheduler/handler/backup_helper.go | 177 ++++++++++++--------- pkg/scheduler/handler/forecast_helper.go | 15 +- pkg/scheduler/handler/handler.go | 23 +-- pkg/scheduler/manager/scheduler.go | 4 +- pkg/scheduler/manager/task_history.go | 4 + pkg/scheduler/models/pending_usage.go | 3 + 24 files changed, 583 insertions(+), 364 deletions(-) create mode 100644 pkg/compute/tasks/ha_disk_create_task.go create mode 100644 pkg/compute/tasks/ha_guest_deploy_task.go create mode 100644 pkg/compute/tasks/ha_guest_start_task.go create mode 100644 pkg/compute/tasks/ha_guest_stop_task.go diff --git a/pkg/apis/compute/guest_const.go b/pkg/apis/compute/guest_const.go index 862e0a20c2..d88531d250 100644 --- a/pkg/apis/compute/guest_const.go +++ b/pkg/apis/compute/guest_const.go @@ -150,7 +150,7 @@ const ( HYPERVISOR_DEFAULT = HYPERVISOR_KVM ) -var VM_RUNNING_STATUS = []string{VM_START_START, VM_STARTING, VM_RUNNING, VM_BLOCK_STREAM} +var VM_RUNNING_STATUS = []string{VM_START_START, VM_STARTING, VM_RUNNING, VM_BLOCK_STREAM, VM_BLOCK_STREAM_FAIL} var VM_CREATING_STATUS = []string{VM_CREATE_NETWORK, VM_CREATE_DISK, VM_START_DEPLOY, VM_DEPLOYING} var HYPERVISORS = []string{ diff --git a/pkg/compute/guestdrivers/virtualization.go b/pkg/compute/guestdrivers/virtualization.go index 817d049c82..f6d202ae1f 100644 --- a/pkg/compute/guestdrivers/virtualization.go +++ b/pkg/compute/guestdrivers/virtualization.go @@ -21,6 +21,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/util/netutils" api "yunion.io/x/onecloud/pkg/apis/compute" @@ -75,6 +76,19 @@ func (self *SVirtualizedGuestDriver) GetRandomNetworkTypes() []string { return []string{api.NETWORK_TYPE_GUEST} } +func (self *SVirtualizedGuestDriver) wireAvaiableForGuest(guest *models.SGuest, wire *models.SWire) (bool, error) { + if guest.BackupHostId == "" { + return true, nil + } else { + backupHost := models.HostManager.FetchHostById(guest.BackupHostId) + count, err := backupHost.GetWiresQuery().Equals("wire_id", wire.Id).CountWithError() + if err != nil { + return false, errors.Wrap(err, "query host wire") + } + return count > 0, nil + } +} + func (self *SVirtualizedGuestDriver) Attach2RandomNetwork(guest *models.SGuest, ctx context.Context, userCred mcclient.TokenCredential, host *models.SHost, netConfig *api.NetworkConfig, pendingUsage quotas.IQuota) ([]models.SGuestnetwork, error) { var wirePattern *regexp.Regexp if len(netConfig.Wire) > 0 { @@ -89,6 +103,11 @@ func (self *SVirtualizedGuestDriver) Attach2RandomNetwork(guest *models.SGuest, for i := 0; i < len(hostwires); i += 1 { hostwire := hostwires[i] wire := hostwire.GetWire() + if ok, err := self.wireAvaiableForGuest(guest, wire); err != nil { + return nil, err + } else if !ok { + continue + } if wire == nil { log.Errorf("host wire is nil?????") @@ -167,7 +186,11 @@ func (self *SVirtualizedGuestDriver) RequestGuestCreateInsertIso(ctx context.Con } func (self *SVirtualizedGuestDriver) StartGuestStopTask(guest *models.SGuest, ctx context.Context, userCred mcclient.TokenCredential, params *jsonutils.JSONDict, parentTaskId string) error { - task, err := taskman.TaskManager.NewTask(ctx, "GuestStopTask", guest, userCred, params, parentTaskId, "", nil) + taskName := "GuestStopTask" + if guest.BackupHostId != "" { + taskName = "HAGuestStopTask" + } + task, err := taskman.TaskManager.NewTask(ctx, taskName, guest, userCred, params, parentTaskId, "", nil) if err != nil { return err } diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index 055a7bad70..21b744ac91 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -559,7 +559,12 @@ func (self *SDisk) StartDiskCreateTask(ctx context.Context, userCred mcclient.To if len(snapshot) > 0 { kwargs.Add(jsonutils.NewString(snapshot), "snapshot") } - if task, err := taskman.TaskManager.NewTask(ctx, "DiskCreateTask", self, userCred, kwargs, parentTaskId, "", nil); err != nil { + + taskName := "DiskCreateTask" + if self.BackupStorageId != "" { + taskName = "HADiskCreateTask" + } + if task, err := taskman.TaskManager.NewTask(ctx, taskName, self, userCred, kwargs, parentTaskId, "", nil); err != nil { return err } else { task.ScheduleRun(nil) diff --git a/pkg/compute/models/guest_actions.go b/pkg/compute/models/guest_actions.go index 4983c35b85..27cb462c59 100644 --- a/pkg/compute/models/guest_actions.go +++ b/pkg/compute/models/guest_actions.go @@ -798,13 +798,21 @@ func (self *SGuest) PerformStart(ctx context.Context, userCred mcclient.TokenCre } } -func (self *SGuest) StartGuestDeployTask(ctx context.Context, userCred mcclient.TokenCredential, kwargs *jsonutils.JSONDict, action string, parentTaskId string) error { +func (self *SGuest) StartGuestDeployTask( + ctx context.Context, userCred mcclient.TokenCredential, + kwargs *jsonutils.JSONDict, action string, parentTaskId string, +) error { self.SetStatus(userCred, api.VM_START_DEPLOY, "") if kwargs == nil { kwargs = jsonutils.NewDict() } kwargs.Add(jsonutils.NewString(action), "deploy_action") - task, err := taskman.TaskManager.NewTask(ctx, "GuestDeployTask", self, userCred, kwargs, parentTaskId, "", nil) + + taskName := "GuestDeployTask" + if self.BackupHostId != "" { + taskName = "HAGuestDeployTask" + } + task, err := taskman.TaskManager.NewTask(ctx, taskName, self, userCred, kwargs, parentTaskId, "", nil) if err != nil { return err } @@ -1019,7 +1027,11 @@ func (self *SGuest) GuestNonSchedStartTask( data *jsonutils.JSONDict, parentTaskId string, ) error { self.SetStatus(userCred, api.VM_START_START, "") - task, err := taskman.TaskManager.NewTask(ctx, "GuestStartTask", self, userCred, data, parentTaskId, "", nil) + taskName := "GuestStartTask" + if self.BackupHostId != "" { + taskName = "HAGuestStartTask" + } + task, err := taskman.TaskManager.NewTask(ctx, taskName, self, userCred, data, parentTaskId, "", nil) if err != nil { return err } @@ -3168,6 +3180,34 @@ func (self *SGuest) PerformBlockStreamFailed(ctx context.Context, userCred mccli return nil, nil } +func (self *SGuest) AllowPerformSlaveStarted(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { + return db.IsAdminAllowPerform(userCred, self, "slave-started") +} + +func (self *SGuest) PerformSlaveStarted(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { + if self.GetMetadata("__mirror_job_status", userCred) == "failed" { + if port, err := data.Int("nbd_server_port"); err != nil { + return nil, httperrors.NewMissingParameterError("nbd_server_port") + } else { + self.StartMirrorJob(ctx, userCred, port, "") + } + } + return nil, nil +} + +func (self *SGuest) StartMirrorJob(ctx context.Context, userCred mcclient.TokenCredential, nbdServerPort int64, parentTaskId string) error { + taskData := jsonutils.NewDict() + taskData.Set("nbd_server_port", jsonutils.NewInt(nbdServerPort)) + if task, err := taskman.TaskManager.NewTask( + ctx, "GuestReSyncToBackup", self, userCred, taskData, parentTaskId, "", nil); err != nil { + log.Errorln(err) + return err + } else { + task.ScheduleRun(nil) + } + return nil +} + func (manager *SGuestManager) AllowPerformDirtyServerStart(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool { return db.IsAdminAllowClassPerform(userCred, manager, "dirty-server-start") } diff --git a/pkg/compute/tasks/disk_create_task.go b/pkg/compute/tasks/disk_create_task.go index a5f49da85f..5d56c5c502 100644 --- a/pkg/compute/tasks/disk_create_task.go +++ b/pkg/compute/tasks/disk_create_task.go @@ -39,25 +39,13 @@ func (self *DiskCreateTask) OnInit(ctx context.Context, obj db.IStandaloneModel, storagecache := disk.GetStorage().GetStoragecache() imageId := disk.GetTemplateId() if len(imageId) > 0 { - if len(disk.BackupStorageId) > 0 { - self.SetStage("OnMasterStorageCacheImageComplete", nil) - } else { - self.SetStage("OnStorageCacheImageComplete", nil) - } + self.SetStage("OnStorageCacheImageComplete", nil) storagecache.StartImageCacheTask(ctx, self.UserCred, imageId, disk.DiskFormat, false, self.GetTaskId()) } else { self.OnStorageCacheImageComplete(ctx, disk, nil) } } -func (self *DiskCreateTask) OnMasterStorageCacheImageComplete(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) { - storage := models.StorageManager.FetchStorageById(disk.BackupStorageId) - storagecache := storage.GetStoragecache() - imageId := disk.GetTemplateId() - self.SetStage("OnStorageCacheImageComplete", nil) - storagecache.StartImageCacheTask(ctx, self.UserCred, imageId, disk.DiskFormat, false, self.GetTaskId()) -} - func (self *DiskCreateTask) OnStorageCacheImageComplete(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) { rebuild, _ := self.GetParams().Bool("rebuild") snapshot, _ := self.GetParams().GetString("snapshot") @@ -68,41 +56,18 @@ func (self *DiskCreateTask) OnStorageCacheImageComplete(ctx context.Context, dis host := storage.GetMasterHost() db.OpsLog.LogEvent(disk, db.ACT_ALLOCATING, disk.GetShortDesc(ctx), self.GetUserCred()) disk.SetStatus(self.GetUserCred(), api.DISK_STARTALLOC, fmt.Sprintf("Disk start alloc use host %s(%s)", host.Name, host.Id)) - if len(disk.BackupStorageId) > 0 { - self.SetStage("OnMasterStorageCreateDiskComplete", nil) - } else { - if rebuild && storage.StorageType == api.STORAGE_RBD { - if count, _ := disk.GetSnapshotCount(); count > 0 { - backingDiskId := stringutils.UUID4() - self.Params.Set("backing_disk_id", jsonutils.NewString(backingDiskId)) - } + if rebuild && storage.StorageType == api.STORAGE_RBD { + if count, _ := disk.GetSnapshotCount(); count > 0 { + backingDiskId := stringutils.UUID4() + self.Params.Set("backing_disk_id", jsonutils.NewString(backingDiskId)) } - self.SetStage("OnDiskReady", nil) } + self.SetStage("OnDiskReady", nil) if err := disk.StartAllocate(ctx, host, storage, self.GetTaskId(), self.GetUserCred(), rebuild, snapshot, self); err != nil { self.OnStartAllocateFailed(ctx, disk, jsonutils.NewString(err.Error())) } } -func (self *DiskCreateTask) OnMasterStorageCreateDiskComplete(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) { - rebuild, _ := self.GetParams().Bool("rebuild") - snapshot, _ := self.GetParams().GetString("snapshot") - storage := models.StorageManager.FetchStorageById(disk.BackupStorageId) - host := storage.GetMasterHost() - db.OpsLog.LogEvent(disk, db.ACT_BACKUP_ALLOCATING, disk.GetShortDesc(ctx), self.GetUserCred()) - disk.SetStatus(self.UserCred, api.DISK_BACKUP_STARTALLOC, fmt.Sprintf("Backup disk start alloc use host %s(%s)", host.Name, host.Id)) - self.SetStage("OnDiskReady", nil) - if err := disk.StartAllocate(ctx, host, storage, self.GetTaskId(), self.GetUserCred(), rebuild, snapshot, self); err != nil { - self.OnBackupAllocateFailed(ctx, disk, jsonutils.NewString(fmt.Sprintf("Backup disk alloctate failed: %s", err.Error()))) - } -} - -func (self *DiskCreateTask) OnBackupAllocateFailed(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) { - disk.SetStatus(self.UserCred, api.DISK_BACKUP_ALLOC_FAILED, data.String()) - logclient.AddActionLogWithStartable(self, disk, logclient.ACT_ALLOCATE, data, self.UserCred, false) - self.SetStageFailed(ctx, data.String()) -} - func (self *DiskCreateTask) OnStartAllocateFailed(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) { disk.SetStatus(self.UserCred, api.DISK_ALLOC_FAILED, data.String()) logclient.AddActionLogWithStartable(self, disk, logclient.ACT_ALLOCATE, data, self.UserCred, false) @@ -144,32 +109,6 @@ func (self *DiskCreateTask) OnDiskReadyFailed(ctx context.Context, disk *models. self.SetStageFailed(ctx, data.String()) } -type DiskCreateBackupTask struct { - DiskCreateTask -} - -func (self *DiskCreateBackupTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { - disk := obj.(*models.SDisk) - storage := models.StorageManager.FetchStorageById(disk.BackupStorageId) - storagecache := storage.GetStoragecache() - imageId := disk.GetTemplateId() - if len(imageId) > 0 { - self.SetStage("OnMasterStorageCreateDiskComplete", nil) - storagecache.StartImageCacheTask(ctx, self.UserCred, imageId, disk.DiskFormat, false, self.GetTaskId()) - } else { - self.OnMasterStorageCreateDiskComplete(ctx, disk, nil) - } -} - -func (self *DiskCreateBackupTask) OnDiskReady(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) { - bkStorage := models.StorageManager.FetchStorageById(disk.BackupStorageId) - bkStorage.ClearSchedDescCache() - disk.SetStatus(self.UserCred, api.DISK_READY, "") - db.OpsLog.LogEvent(disk, db.ACT_BACKUP_ALLOCATE, disk.GetShortDesc(ctx), self.UserCred) - self.SetStageComplete(ctx, nil) -} - func init() { taskman.RegisterTask(DiskCreateTask{}) - taskman.RegisterTask(DiskCreateBackupTask{}) } diff --git a/pkg/compute/tasks/guest_backup_tasks.go b/pkg/compute/tasks/guest_backup_tasks.go index 9f75dc22fd..157a76396f 100644 --- a/pkg/compute/tasks/guest_backup_tasks.go +++ b/pkg/compute/tasks/guest_backup_tasks.go @@ -193,6 +193,7 @@ func (self *GuestStartAndSyncToBackupTask) OnStartBackupGuestFailed(ctx context. } func (self *GuestStartAndSyncToBackupTask) OnRequestSyncToBackup(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + guest.SetMetadata(ctx, "__mirror_job_status", "", self.UserCred) guest.SetStatus(self.UserCred, api.VM_BLOCK_STREAM, "OnSyncToBackup") self.SetStageComplete(ctx, nil) } @@ -350,4 +351,21 @@ func init() { taskman.RegisterTask(GuestSwitchToBackupTask{}) taskman.RegisterTask(GuestStartAndSyncToBackupTask{}) taskman.RegisterTask(GuestCreateBackupTask{}) + taskman.RegisterTask(GuestReSyncToBackup{}) +} + +type GuestReSyncToBackup struct { + GuestStartAndSyncToBackupTask +} + +func (self *GuestReSyncToBackup) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + guest := obj.(*models.SGuest) + self.StartSyncToBackup(ctx, guest) +} + +func (self *GuestReSyncToBackup) StartSyncToBackup(ctx context.Context, guest *models.SGuest) { + data := jsonutils.NewDict() + nbdServerPort, _ := self.Params.Int("nbd_server_port") + data.Set("nbd_server_port", jsonutils.NewInt(nbdServerPort)) + self.OnStartBackupGuest(ctx, guest, data) } diff --git a/pkg/compute/tasks/guest_deploy_task.go b/pkg/compute/tasks/guest_deploy_task.go index fbc1498e64..2371fb642d 100644 --- a/pkg/compute/tasks/guest_deploy_task.go +++ b/pkg/compute/tasks/guest_deploy_task.go @@ -54,29 +54,12 @@ func (self *GuestDeployTask) OnDeployWaitServerStop(ctx context.Context, guest * self.SetStage("OnDeployGuestComplete", nil) targetHostId, _ := self.Params.GetString("target_host_id") if len(targetHostId) == 0 { - if len(guest.BackupHostId) > 0 { - self.SetStage("OnSlaveHostDeployComplete", nil) - self.DeployBackup(ctx, guest, nil) - return - } else { - targetHostId = guest.HostId - } + targetHostId = guest.HostId } host := models.HostManager.FetchHostById(targetHostId) self.DeployOnHost(ctx, guest, host) } -func (self *GuestDeployTask) DeployBackup(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - host := models.HostManager.FetchHostById(guest.BackupHostId) - err := guest.GetDriver().RequestDeployGuestOnHost(ctx, guest, host, self) - if err != nil { - log.Errorf("request_deploy_guest_on_host %s", err) - self.OnDeployGuestFail(ctx, guest, err) - } else { - guest.SetStatus(self.UserCred, api.VM_DEPLOYING_BACKUP, "") - } -} - func (self *GuestDeployTask) DeployOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost) { err := guest.GetDriver().RequestDeployGuestOnHost(ctx, guest, host, self) if err != nil { @@ -87,16 +70,6 @@ func (self *GuestDeployTask) DeployOnHost(ctx context.Context, guest *models.SGu } } -func (self *GuestDeployTask) OnSlaveHostDeployComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - host := guest.GetHost() - self.SetStage("OnDeployGuestComplete", nil) - self.DeployOnHost(ctx, guest, host) -} - -func (self *GuestDeployTask) OnSlaveHostDeployCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - self.OnDeployGuestFail(ctx, guest, fmt.Errorf("deploy backup failed %s", data)) -} - func (self *GuestDeployTask) OnDeployGuestFail(ctx context.Context, guest *models.SGuest, err error) { guest.SetStatus(self.UserCred, api.VM_DEPLOY_FAILED, err.Error()) self.SetStageFailed(ctx, err.Error()) @@ -182,29 +155,6 @@ func (self *GuestDeployTask) OnDeployGuestSyncstatusCompleteFailed(ctx context.C self.SetStageFailed(ctx, data.String()) } -type GuestDeployBackupTask struct { - GuestDeployTask -} - -func (self *GuestDeployBackupTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { - guest := obj.(*models.SGuest) - if len(guest.BackupHostId) == 0 { - self.SetStageFailed(ctx, "Guest dosen't have backup host") - } - self.SetStage("OnDeployGuestComplete", nil) - self.DeployBackup(ctx, guest, nil) -} - -func (self *GuestDeployBackupTask) OnDeployGuestComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - self.SetStageComplete(ctx, nil) -} - -func (self *GuestDeployBackupTask) OnDeployGuestCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - guest.SetStatus(self.UserCred, api.VM_DEPLOYING_BACKUP_FAILED, data.String()) - self.SetStageComplete(ctx, nil) -} - func init() { taskman.RegisterTask(GuestDeployTask{}) - taskman.RegisterTask(GuestDeployBackupTask{}) } diff --git a/pkg/compute/tasks/guest_start_task.go b/pkg/compute/tasks/guest_start_task.go index a6504a262a..51ceda7a2a 100644 --- a/pkg/compute/tasks/guest_start_task.go +++ b/pkg/compute/tasks/guest_start_task.go @@ -38,46 +38,8 @@ func init() { func (self *GuestStartTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { guest := obj.(*models.SGuest) - self.checkTemplate(ctx, guest) -} - -func (self *GuestStartTask) checkTemplate(ctx context.Context, guest *models.SGuest) { - /*diskCat := guest.CategorizeDisks() - if diskCat.Root != nil && len(diskCat.Root.GetTemplateId()) > 0 { - if len(guest.BackupHostId) > 0 { - self.SetStage("OnMasterHostTemplateReady", nil) - } else { - self.SetStage("OnStartTemplateReady", nil) - } - guest.GetDriver().CheckDiskTemplateOnStorage(ctx, self.UserCred, diskCat.Root.GetTemplateId(), diskCat.Root.DiskFormat, diskCat.Root.StorageId, self) - } else { - self.startStart(ctx, guest) - }*/ - self.startStart(ctx, guest) -} - -func (self *GuestStartTask) OnMasterHostTemplateReady(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - self.SetStage("OnStartTemplateReady", nil) - diskCat := guest.CategorizeDisks() - err := guest.GetDriver().CheckDiskTemplateOnStorage(ctx, self.UserCred, diskCat.Root.GetTemplateId(), diskCat.Root.DiskFormat, - diskCat.Root.BackupStorageId, self) - if err != nil { - self.SetStageFailed(ctx, err.Error()) - } -} - -func (self *GuestStartTask) OnStartTemplateReady(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { - guest := obj.(*models.SGuest) - self.startStart(ctx, guest) -} - -func (self *GuestStartTask) startStart(ctx context.Context, guest *models.SGuest) { db.OpsLog.LogEvent(guest, db.ACT_STARTING, nil, self.UserCred) - if len(guest.BackupHostId) > 0 { - self.RequestStartBacking(ctx, guest) - } else { - self.RequestStart(ctx, guest) - } + self.RequestStart(ctx, guest) } func (self *GuestStartTask) RequestStart(ctx context.Context, guest *models.SGuest) { @@ -96,39 +58,6 @@ func (self *GuestStartTask) RequestStart(ctx context.Context, guest *models.SGue } } -func (self *GuestStartTask) RequestStartBacking(ctx context.Context, guest *models.SGuest) { - self.SetStage("OnStartBackupGuestComplete", nil) - host := models.HostManager.FetchHostById(guest.BackupHostId) - guest.SetStatus(self.UserCred, api.VM_BACKUP_STARTING, "") - result, err := guest.GetDriver().RequestStartOnHost(ctx, guest, host, self.UserCred, self) - if err != nil { - self.OnStartCompleteFailed(ctx, guest, jsonutils.NewString(err.Error())) - } else { - if result != nil && jsonutils.QueryBoolean(result, "is_running", false) { - self.OnStartBackupGuestComplete(ctx, guest, nil) - } - } -} - -func (self *GuestStartTask) OnStartBackupGuestComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - if data != nil { - nbdServerPort, err := data.Int("nbd_server_port") - if err == nil { - backupHost := models.HostManager.FetchHostById(guest.BackupHostId) - nbdServerUri := fmt.Sprintf("nbd:%s:%d", backupHost.AccessIp, nbdServerPort) - guest.SetMetadata(ctx, "backup_nbd_server_uri", nbdServerUri, self.UserCred) - } else { - self.OnStartCompleteFailed(ctx, guest, jsonutils.NewString("Start backup guest result missing nbd_server_port")) - return - } - } - self.RequestStart(ctx, guest) -} - -func (self *GuestStartTask) OnStartBackupGuestCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - self.OnStartCompleteFailed(ctx, guest, data) -} - func (self *GuestStartTask) OnStartComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { guest := obj.(*models.SGuest) db.OpsLog.LogEvent(guest, db.ACT_START, guest.GetShortDesc(ctx), self.UserCred) diff --git a/pkg/compute/tasks/guest_stop_task.go b/pkg/compute/tasks/guest_stop_task.go index f88e31bf77..0ee7ffef1a 100644 --- a/pkg/compute/tasks/guest_stop_task.go +++ b/pkg/compute/tasks/guest_stop_task.go @@ -50,7 +50,7 @@ func (self *GuestStopTask) stopGuest(ctx context.Context, guest *models.SGuest) if !self.IsSubtask() { guest.SetStatus(self.UserCred, api.VM_STOPPING, "") } - self.SetStage("OnMasterStopTaskComplete", nil) + self.SetStage("OnGuestStopTaskComplete", nil) err := guest.GetDriver().RequestStopOnHost(ctx, guest, host, self) if err != nil { log.Errorf("RequestStopOnHost fail %s", err) @@ -58,25 +58,6 @@ func (self *GuestStopTask) stopGuest(ctx context.Context, guest *models.SGuest) } } -func (self *GuestStopTask) OnMasterStopTaskComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { - if len(guest.BackupHostId) > 0 { - host := models.HostManager.FetchHostById(guest.BackupHostId) - self.SetStage("OnGuestStopTaskComplete", nil) - err := guest.GetDriver().RequestStopOnHost(ctx, guest, host, self) - if err != nil { - log.Errorf("RequestStopOnHost fail %s", err) - self.OnGuestStopTaskCompleteFailed(ctx, guest, jsonutils.NewString(err.Error())) - } - } else { - self.OnGuestStopTaskComplete(ctx, guest, data) - } -} - -func (self *GuestStopTask) OnMasterStopTaskCompleteFailed(ctx context.Context, obj db.IStandaloneModel, reason jsonutils.JSONObject) { - guest := obj.(*models.SGuest) - self.OnGuestStopTaskCompleteFailed(ctx, guest, reason) -} - func (self *GuestStopTask) OnGuestStopTaskComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { if !self.IsSubtask() { guest.StartSyncstatus(ctx, self.UserCred, "") diff --git a/pkg/compute/tasks/ha_disk_create_task.go b/pkg/compute/tasks/ha_disk_create_task.go new file mode 100644 index 0000000000..b465347e67 --- /dev/null +++ b/pkg/compute/tasks/ha_disk_create_task.go @@ -0,0 +1,100 @@ +package tasks + +import ( + "context" + "fmt" + + "yunion.io/x/jsonutils" + + api "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" +) + +func init() { + taskman.RegisterTask(HADiskCreateTask{}) + taskman.RegisterTask(DiskCreateBackupTask{}) +} + +type HADiskCreateTask struct { + DiskCreateTask +} + +func (self *HADiskCreateTask) OnStorageCacheImageComplete( + ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject, +) { + storage := models.StorageManager.FetchStorageById(disk.BackupStorageId) + storagecache := storage.GetStoragecache() + imageId := disk.GetTemplateId() + self.SetStage("OnBackupStorageCacheImageComplete", nil) + storagecache.StartImageCacheTask(ctx, self.UserCred, imageId, disk.DiskFormat, false, self.GetTaskId()) +} + +func (self *HADiskCreateTask) OnBackupStorageCacheImageComplete( + ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject, +) { + self.DiskCreateTask.OnStorageCacheImageComplete(ctx, disk, data) +} + +func (self *HADiskCreateTask) OnDiskReady( + ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject, +) { + rebuild, _ := self.GetParams().Bool("rebuild") + snapshot, _ := self.GetParams().GetString("snapshot") + storage := models.StorageManager.FetchStorageById(disk.BackupStorageId) + host := storage.GetMasterHost() + db.OpsLog.LogEvent(disk, db.ACT_BACKUP_ALLOCATING, disk.GetShortDesc(ctx), self.GetUserCred()) + disk.SetStatus(self.UserCred, api.DISK_BACKUP_STARTALLOC, + fmt.Sprintf("Backup disk start alloc use host %s(%s)", host.Name, host.Id), + ) + self.SetStage("OnSlaveDiskReady", nil) + if err := disk.StartAllocate(ctx, host, storage, + self.GetTaskId(), self.GetUserCred(), rebuild, snapshot, self, + ); err != nil { + self.OnBackupAllocateFailed(ctx, disk, + jsonutils.NewString(fmt.Sprintf("Backup disk alloctate failed: %s", err.Error()))) + } +} + +func (self *HADiskCreateTask) OnBackupAllocateFailed(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) { + disk.SetStatus(self.UserCred, api.DISK_BACKUP_ALLOC_FAILED, data.String()) + self.SetStageFailed(ctx, data.String()) +} + +func (self *HADiskCreateTask) OnSlaveDiskReady( + ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject, +) { + self.DiskCreateTask.OnDiskReady(ctx, disk, data) +} + +type DiskCreateBackupTask struct { + HADiskCreateTask +} + +func (self *DiskCreateBackupTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + disk := obj.(*models.SDisk) + storage := models.StorageManager.FetchStorageById(disk.BackupStorageId) + storagecache := storage.GetStoragecache() + imageId := disk.GetTemplateId() + if len(imageId) > 0 { + self.SetStage("OnBackupStorageCacheImageComplete", nil) + storagecache.StartImageCacheTask(ctx, self.UserCred, imageId, disk.DiskFormat, false, self.GetTaskId()) + } else { + self.OnBackupStorageCacheImageComplete(ctx, disk, nil) + } +} + +func (self *DiskCreateBackupTask) OnBackupStorageCacheImageComplete( + ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject, +) { + self.HADiskCreateTask.OnDiskReady(ctx, disk, data) +} + +func (self *DiskCreateBackupTask) OnSlaveDiskReady(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) { + bkStorage := models.StorageManager.FetchStorageById(disk.BackupStorageId) + bkStorage.ClearSchedDescCache() + disk.SetStatus(self.UserCred, api.DISK_READY, "") + db.OpsLog.LogEvent(disk, db.ACT_BACKUP_ALLOCATE, disk.GetShortDesc(ctx), self.UserCred) + self.SetStageComplete(ctx, nil) +} diff --git a/pkg/compute/tasks/ha_guest_deploy_task.go b/pkg/compute/tasks/ha_guest_deploy_task.go new file mode 100644 index 0000000000..cc1d27e8dd --- /dev/null +++ b/pkg/compute/tasks/ha_guest_deploy_task.go @@ -0,0 +1,74 @@ +package tasks + +import ( + "context" + "fmt" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + + api "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" +) + +func init() { + taskman.RegisterTask(GuestDeployBackupTask{}) + taskman.RegisterTask(HAGuestDeployTask{}) +} + +type HAGuestDeployTask struct { + GuestDeployTask +} + +func (self *HAGuestDeployTask) OnDeployGuestComplete( + ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, +) { + self.DeployBackup(ctx, guest, data) +} + +func (self *HAGuestDeployTask) DeployBackup(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + host := models.HostManager.FetchHostById(guest.BackupHostId) + err := guest.GetDriver().RequestDeployGuestOnHost(ctx, guest, host, self) + if err != nil { + log.Errorf("request_deploy_guest_on_host %s", err) + self.OnDeployGuestFail(ctx, guest, err) + } else { + guest.SetStatus(self.UserCred, api.VM_DEPLOYING_BACKUP, "") + } +} + +func (self *HAGuestDeployTask) OnDeploySlaveGuestComplete( + ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, +) { + self.GuestDeployTask.OnDeployGuestComplete(ctx, guest, data) +} + +func (self *HAGuestDeployTask) OnDeploySlaveGuestCompleteFailed( + ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, +) { + self.OnDeployGuestFail(ctx, guest, fmt.Errorf("deploy backup failed %s", data)) +} + +type GuestDeployBackupTask struct { + HAGuestDeployTask +} + +func (self *GuestDeployBackupTask) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { + guest := obj.(*models.SGuest) + if len(guest.BackupHostId) == 0 { + self.SetStageFailed(ctx, "Guest dosen't have backup host") + } + self.SetStage("OnDeployGuestComplete", nil) + self.DeployBackup(ctx, guest, nil) +} + +func (self *GuestDeployBackupTask) OnDeployGuestComplete(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + self.SetStageComplete(ctx, nil) +} + +func (self *GuestDeployBackupTask) OnDeployGuestCompleteFailed(ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject) { + 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 new file mode 100644 index 0000000000..87b2203943 --- /dev/null +++ b/pkg/compute/tasks/ha_guest_start_task.go @@ -0,0 +1,67 @@ +package tasks + +import ( + "context" + "fmt" + + "yunion.io/x/jsonutils" + + api "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" +) + +func init() { + taskman.RegisterTask(HAGuestStartTask{}) +} + +type HAGuestStartTask struct { + GuestStartTask +} + +func (self *HAGuestStartTask) OnInit( + ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject, +) { + guest := obj.(*models.SGuest) + db.OpsLog.LogEvent(guest, db.ACT_STARTING, nil, self.UserCred) + self.RequestStartBacking(ctx, guest) +} + +func (self *HAGuestStartTask) RequestStartBacking(ctx context.Context, guest *models.SGuest) { + self.SetStage("OnStartBackupGuestComplete", nil) + host := models.HostManager.FetchHostById(guest.BackupHostId) + guest.SetStatus(self.UserCred, api.VM_BACKUP_STARTING, "") + result, err := guest.GetDriver().RequestStartOnHost(ctx, guest, host, self.UserCred, self) + if err != nil { + self.OnStartCompleteFailed(ctx, guest, jsonutils.NewString(err.Error())) + } else { + if result != nil && jsonutils.QueryBoolean(result, "is_running", false) { + self.OnStartBackupGuestComplete(ctx, guest, nil) + } + } +} + +func (self *HAGuestStartTask) OnStartBackupGuestComplete( + ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, +) { + if data != nil { + nbdServerPort, err := data.Int("nbd_server_port") + if err == nil { + backupHost := models.HostManager.FetchHostById(guest.BackupHostId) + nbdServerUri := fmt.Sprintf("nbd:%s:%d", backupHost.AccessIp, nbdServerPort) + guest.SetMetadata(ctx, "backup_nbd_server_uri", nbdServerUri, self.UserCred) + } else { + self.OnStartCompleteFailed(ctx, guest, + jsonutils.NewString("Start backup guest result missing nbd_server_port")) + return + } + } + self.RequestStart(ctx, guest) +} + +func (self *HAGuestStartTask) OnStartBackupGuestCompleteFailed( + ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, +) { + 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 new file mode 100644 index 0000000000..b1630bda23 --- /dev/null +++ b/pkg/compute/tasks/ha_guest_stop_task.go @@ -0,0 +1,44 @@ +package tasks + +import ( + "context" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" + "yunion.io/x/onecloud/pkg/compute/models" +) + +func init() { + taskman.RegisterTask(HAGuestStopTask{}) +} + +type HAGuestStopTask struct { + GuestStopTask +} + +func (self *HAGuestStopTask) OnGuestStopTaskComplete( + ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, +) { + host := models.HostManager.FetchHostById(guest.BackupHostId) + self.SetStage("OnSlaveGuestStopTaskComplete", nil) + err := guest.GetDriver().RequestStopOnHost(ctx, guest, host, self) + if err != nil { + log.Errorf("RequestStopOnHost fail %s", err) + self.OnGuestStopTaskCompleteFailed( + ctx, guest, jsonutils.NewString(err.Error())) + } +} + +func (self *HAGuestStopTask) OnSlaveGuestStopTaskComplete( + ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, +) { + self.GuestStopTask.OnGuestStopTaskComplete(ctx, guest, data) +} + +func (self *HAGuestStopTask) OnSlaveGuestStopTaskCompleteFailed( + ctx context.Context, guest *models.SGuest, data jsonutils.JSONObject, +) { + self.OnGuestStopTaskCompleteFailed(ctx, guest, data) +} diff --git a/pkg/hostman/guestman/guesttasks.go b/pkg/hostman/guestman/guesttasks.go index 0951733135..441948cb86 100644 --- a/pkg/hostman/guestman/guesttasks.go +++ b/pkg/hostman/guestman/guesttasks.go @@ -978,9 +978,6 @@ func (s *SDriveMirrorTask) startMirror(res string) { } disks, _ := s.Desc.GetArray("disks") if s.index < len(disks) { - if s.index >= 1 { // data disk - s.syncMode = "none" - } target := fmt.Sprintf("%s:exportname=drive_%d", s.nbdUri, s.index) s.Monitor.DriveMirror(s.startMirror, fmt.Sprintf("drive_%d", s.index), target, s.syncMode, true) diff --git a/pkg/hostman/guestman/qemu-kvm.go b/pkg/hostman/guestman/qemu-kvm.go index a14249813a..64e975815f 100644 --- a/pkg/hostman/guestman/qemu-kvm.go +++ b/pkg/hostman/guestman/qemu-kvm.go @@ -510,6 +510,18 @@ func (s *SKVMGuestInstance) SyncMirrorJobFailed(reason string) { } } +func (s *SKVMGuestInstance) OnSlaveStartedWithNbdServer(nbdServerPort int64) { + params := jsonutils.NewDict() + params.Set("nbd_server_port", jsonutils.NewInt(nbdServerPort)) + _, err := modules.Servers.PerformAction( + hostutils.GetComputeSession(context.Background()), + s.GetId(), "slave-started", params, + ) + if err != nil { + log.Errorf("Server %s perform resume guest got error %s", s.GetId(), err) + } +} + func (s *SKVMGuestInstance) onMonitorConnected(ctx context.Context) { log.Infof("Monitor connected ...") s.Monitor.GetVersion(func(v string) { @@ -525,11 +537,9 @@ func (s *SKVMGuestInstance) onGetQemuVersion(ctx context.Context, version string body := jsonutils.NewDict() body.Set("live_migrate_dest_port", migratePort) hostutils.TaskComplete(ctx, body) - } else if jsonutils.QueryBoolean(s.Desc, "is_slave", false) { - if ctx != nil && len(appctx.AppContextTaskId(ctx)) > 0 { - s.startQemuBuiltInNbdServer(ctx) - } - } else if jsonutils.QueryBoolean(s.Desc, "is_master", false) { + } else if s.IsSlave() { + s.startQemuBuiltInNbdServer(ctx) + } else if s.IsMaster() { s.startDiskBackupMirror(ctx) if ctx != nil && len(appctx.AppContextTaskId(ctx)) > 0 { s.DoResumeTask(ctx) @@ -582,19 +592,19 @@ func (s *SKVMGuestInstance) startDiskBackupMirror(ctx context.Context) { } func (s *SKVMGuestInstance) startQemuBuiltInNbdServer(ctx context.Context) { - if ctx == nil || len(appctx.AppContextTaskId(ctx)) == 0 { - return - } - nbdServerPort := s.manager.GetFreePortByBase(BUILT_IN_NBD_SERVER_PORT_BASE) var onNbdServerStarted = func(res string) { - if len(res) > 0 { - log.Errorf("Start Qemu Builtin nbd server error %s", res) - hostutils.TaskFailed(ctx, res) + if ctx != nil && len(appctx.AppContextTaskId(ctx)) > 0 { + if len(res) > 0 { + log.Errorf("Start Qemu Builtin nbd server error %s", res) + hostutils.TaskFailed(ctx, res) + } else { + res := jsonutils.NewDict() + res.Set("nbd_server_port", jsonutils.NewInt(int64(nbdServerPort))) + hostutils.TaskComplete(ctx, res) + } } else { - res := jsonutils.NewDict() - res.Set("nbd_server_port", jsonutils.NewInt(int64(nbdServerPort))) - hostutils.TaskComplete(ctx, res) + s.OnSlaveStartedWithNbdServer(int64(nbdServerPort)) } } s.Monitor.StartNbdServer(nbdServerPort, true, true, onNbdServerStarted) diff --git a/pkg/scheduler/api/sched.go b/pkg/scheduler/api/sched.go index ee774f4044..7cd0620193 100644 --- a/pkg/scheduler/api/sched.go +++ b/pkg/scheduler/api/sched.go @@ -31,10 +31,11 @@ import ( type SchedInfo struct { *api.ScheduleInput - Tag string `json:"tag"` - Type string `json:"type"` - IsContainer bool `json:"is_container"` - Candidates []string `json:"candidates"` + Tag string `json:"tag"` + Type string `json:"type"` + IsContainer bool `json:"is_container"` + PreferCandidates []string `json:"candidates"` + RequiredCandidates int `json:"required_candidates"` IgnoreFilters map[string]bool `json:"ignore_filters"` IsSuggestion bool `json:"suggestion"` @@ -88,14 +89,25 @@ func FetchSchedInfo(req *http.Request) (*SchedInfo, error) { func NewSchedInfo(input *api.ScheduleInput) *SchedInfo { data := new(SchedInfo) data.ScheduleInput = input + data.RequiredCandidates = 1 if data.Hypervisor == "" || data.Hypervisor == SchedTypeKvm { data.Hypervisor = HostHypervisorForKvm } - candidates := make([]string, 0) - if !data.Backup && data.PreferHost != "" { - candidates = append(candidates, data.PreferHost) + preferCandidates := make([]string, 0) + if data.PreferHost != "" { + preferCandidates = append(preferCandidates, data.PreferHost) + } + + if data.Backup { + if data.PreferBackupHost != "" { + preferCandidates = append(preferCandidates, data.PreferBackupHost) + } + data.RequiredCandidates += 1 + } else { + // make sure prefer backup is null + data.PreferBackupHost = "" } if data.ResourceType == "" { @@ -108,7 +120,7 @@ func NewSchedInfo(input *api.ScheduleInput) *SchedInfo { log.V(4).Warningf("No baremetal_disk_config info found in json, use default baremetal disk config: %#v", defaultConfs) } - data.Candidates = candidates + data.PreferCandidates = preferCandidates data.reviseData() diff --git a/pkg/scheduler/core/context.go b/pkg/scheduler/core/context.go index 383c9a26ea..d92374ffdb 100644 --- a/pkg/scheduler/core/context.go +++ b/pkg/scheduler/core/context.go @@ -22,6 +22,7 @@ import ( "yunion.io/x/log" "yunion.io/x/pkg/tristate" + "yunion.io/x/pkg/utils" "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/compute/models" @@ -448,19 +449,9 @@ func (u *Unit) SchedData() *api.SchedInfo { } func (u *Unit) ShouldExecuteSchedtagFilter(hostId string) bool { - schedData := u.SchedData() - if !schedData.Backup { - if len(schedData.PreferHost) != 0 { - return false - } - return true - } - for _, preferHost := range []string{schedData.PreferHost, schedData.PreferBackupHost} { - if preferHost == hostId { - return false - } - } - return true + return !utils.IsInStringArray( + hostId, []string{u.SchedInfo.PreferHost, u.SchedInfo.PreferBackupHost}, + ) } func (u *Unit) GetHypervisor() string { diff --git a/pkg/scheduler/core/generic_scheduler.go b/pkg/scheduler/core/generic_scheduler.go index 79b48f4845..0367f5987a 100644 --- a/pkg/scheduler/core/generic_scheduler.go +++ b/pkg/scheduler/core/generic_scheduler.go @@ -193,8 +193,8 @@ func newSchedResultByCtx(u *Unit, count int64, c Candidater) *SchedResultItem { return r } -func generateScheduleResult(u *Unit, scs []*SelectedCandidate, cs []Candidater) ([]*SchedResultItem, error) { - results := make([]*SchedResultItem, 0) +func generateScheduleResult(u *Unit, scs []*SelectedCandidate, cs []Candidater) (SchedResultItems, error) { + results := make(SchedResultItems, 0) itemMap := make(map[string]int) for _, it := range scs { @@ -218,10 +218,9 @@ func generateScheduleResult(u *Unit, scs []*SelectedCandidate, cs []Candidater) } suggestionAll := u.SchedInfo.SuggestionAll - if suggestionAll || len(u.SchedData().Candidates) > 0 { + if suggestionAll || len(u.SchedData().PreferCandidates) > 0 { for _, c := range cs { if suggestionLimit <= int64(len(results)) { - break } id := c.IndexKey() @@ -272,21 +271,18 @@ func GetCapacities(u *Unit, id string) (res map[string]int64) { return } -type SchedResultItemList struct { - Unit *Unit - Data []*SchedResultItem +type SchedResultItems []*SchedResultItem + +func (its SchedResultItems) Len() int { + return len(its) } -func (its SchedResultItemList) Len() int { - return len(its.Data) +func (its SchedResultItems) Swap(i, j int) { + its[i], its[j] = its[j], its[i] } -func (its *SchedResultItemList) Swap(i, j int) { - its.Data[i], its.Data[j] = its.Data[j], its.Data[i] -} - -func (its SchedResultItemList) Less(i, j int) bool { - it1, it2 := its.Data[i], its.Data[j] +func (its SchedResultItems) Less(i, j int) bool { + it1, it2 := its[i], its[j] return it1.Capacity < it2.Capacity /* ctx := its.Unit @@ -310,6 +306,11 @@ func (its SchedResultItemList) Less(i, j int) bool { */ } +type SchedResultItemList struct { + Unit *Unit + Data SchedResultItems +} + func (its SchedResultItemList) String() string { bytes, _ := json.Marshal(its.Data) return string(bytes) diff --git a/pkg/scheduler/handler/backup_helper.go b/pkg/scheduler/handler/backup_helper.go index 01c41d6b94..f523dffdf8 100644 --- a/pkg/scheduler/handler/backup_helper.go +++ b/pkg/scheduler/handler/backup_helper.go @@ -22,10 +22,11 @@ import ( schedapi "yunion.io/x/onecloud/pkg/apis/scheduler" "yunion.io/x/onecloud/pkg/scheduler/core" - //schedman "yunion.io/x/onecloud/pkg/scheduler/manager" ) -func transToBackupSchedResult(result *core.SchedResultItemList, preferMasterHost, preferBackupHost string, count int64, sid string) *schedapi.ScheduleOutput { +func transToBackupSchedResult( + result *core.SchedResultItemList, preferMasterHost, preferBackupHost string, count int64, sid string, +) *schedapi.ScheduleOutput { // clean each result sched result item's count for _, item := range result.Data { item.Count = 0 @@ -43,9 +44,10 @@ func newBackupSchedResult( ) *schedapi.ScheduleOutput { ret := new(schedapi.ScheduleOutput) apiResults := make([]*schedapi.CandidateResource, 0) + var wireHostMap map[string]core.SchedResultItems for i := 0; i < int(count); i++ { log.V(10).Debugf("Select backup host from result: %s", result) - target, err := getSchedBackupResult(result, preferMasterHost, preferBackupHost, sid) + target, err := getSchedBackupResult(result, preferMasterHost, preferBackupHost, sid, wireHostMap) if err != nil { er := &schedapi.CandidateResource{Error: err.Error()} apiResults = append(apiResults, er) @@ -60,20 +62,24 @@ func newBackupSchedResult( func getSchedBackupResult( result *core.SchedResultItemList, preferMasterHost, preferBackupHost string, - sid string, + sid string, wireHostMap map[string]core.SchedResultItems, ) (*schedapi.CandidateResource, error) { - masterHost := selectMasterHost(result.Data, preferMasterHost, preferBackupHost) + if wireHostMap == nil { + wireHostMap = buildWireHostMap(result) + } else { + reviseWireHostMap(wireHostMap) + } + + masterHost, backupHost := selectHosts(wireHostMap, preferMasterHost, preferBackupHost) if masterHost == nil { return nil, fmt.Errorf("Can't find master host %q", preferMasterHost) } - backupHost := selectBackupHost(masterHost.ID, preferBackupHost, result.Data) if backupHost == nil { return nil, fmt.Errorf("Can't find backup host %q by master %q", preferBackupHost, masterHost.ID) } markHostUsed(masterHost) markHostUsed(backupHost) - sort.Sort(sort.Reverse(result)) ret := masterHost.ToCandidateResource() ret.BackupCandidate = backupHost.ToCandidateResource() @@ -82,80 +88,109 @@ func getSchedBackupResult( return ret, nil } +func buildWireHostMap(result *core.SchedResultItemList) map[string]core.SchedResultItems { + sort.Sort(sort.Reverse(result.Data)) + wireHostMap := make(map[string]core.SchedResultItems) + for i := 0; i < len(result.Data); i++ { + networks := result.Data[i].Candidater.Getter().Networks() + for j := 0; j < len(networks); j++ { + if hosts, ok := wireHostMap[networks[j].WireId]; ok { + if hostInResultItemsIndex(result.Data[i].ID, hosts) < 0 { + wireHostMap[networks[j].WireId] = append(hosts, result.Data[i]) + } + } else { + wireHostMap[networks[j].WireId] = core.SchedResultItems{result.Data[i]} + } + } + } + return wireHostMap +} + +func reviseWireHostMap(wireHostMap map[string]core.SchedResultItems) { + for _, hosts := range wireHostMap { + sort.Sort(sort.Reverse(hosts)) + } +} + func markHostUsed(host *core.SchedResultItem) { host.Count++ host.Capacity-- } -// selectMasterID find master host id run VM -// return nil if not found -func selectMasterHost(result []*core.SchedResultItem, preferMasterHost, preferBackupHost string) *core.SchedResultItem { - if len(result) == 0 { - return nil - } - host := result[0] - if host.ID == preferMasterHost { - if host.Capacity >= 1 { - return host - } else { - return nil +func hostInResultItemsIndex(hostId string, hosts core.SchedResultItems) int { + for i := 0; i < len(hosts); i++ { + if hosts[i].ID == hostId { + return i } } - if host.Capacity >= 1 && host.ID != preferBackupHost { - if len(preferMasterHost) == 0 { - return host - } - if len(result) == 1 { - return nil - } - restHosts := result[1:] - return selectMasterHost(restHosts, preferMasterHost, preferBackupHost) - } - if len(result) == 1 { - return nil - } - restHosts := result[1:] - return selectMasterHost(restHosts, preferMasterHost, preferBackupHost) + return -1 } -func selectBackupHost(masterID, preferBackupHost string, result []*core.SchedResultItem) *core.SchedResultItem { - if len(result) == 0 { - return nil - } - firstHost := result[0] - if canHostAsBackup(masterID, preferBackupHost, firstHost) { - return firstHost - } - if len(result) == 1 { - return nil - } - restHosts := result[1:] - return selectBackupHost(masterID, preferBackupHost, restHosts) -} +func selectHosts( + wireHostMap map[string]core.SchedResultItems, preferMasterHost, preferBackupHost string, +) (*core.SchedResultItem, *core.SchedResultItem) { + var scroe int64 + var masterIdx, backupIdx int + var selectedWireId string + for wireId, hosts := range wireHostMap { + masterIdx, backupIdx = -1, -1 + if len(hosts) < 2 { + continue + } + if len(preferMasterHost) > 0 { + if masterIdx = hostInResultItemsIndex(preferMasterHost, hosts); masterIdx < 0 { + continue + } + } + if len(preferBackupHost) > 0 { + if backupIdx = hostInResultItemsIndex(preferBackupHost, hosts); backupIdx < 0 { + continue + } + } -func canHostAsBackup(masterID, preferBackupHost string, host *core.SchedResultItem) bool { - if host.ID == masterID { - return false - } - if host.Capacity == 0 { - return false - } - if preferBackupHost != "" { - if host.ID != preferBackupHost { - return false + // select master host index + if masterIdx < 0 { + for i := 0; i < len(hosts); i++ { + if hosts[i].ID != preferBackupHost { + masterIdx = i + } + } + } + if hosts[masterIdx].Capacity <= 0 { + if len(preferMasterHost) > 0 { + // in case prefer master host capacity isn't enough + break + } else { + continue + } + } + + // select backup host index + if backupIdx < 0 { + for i := 0; i < len(hosts); i++ { + if i != masterIdx { + backupIdx = i + } + } + } + if hosts[backupIdx].Capacity <= 0 { + if len(preferBackupHost) > 0 { + // in case perfer backup host capacity isn't enough + break + } else { + continue + } + } + + // the highest total score wins + curScore := hosts[masterIdx].Capacity + hosts[backupIdx].Capacity + if curScore > scroe { + selectedWireId = wireId + scroe = curScore } } - return true -} - -type dirtyItemAdapter struct { - *core.SchedResultItem -} - -func (a *dirtyItemAdapter) Index() (string, error) { - return a.ID, nil -} - -func (a *dirtyItemAdapter) GetCount() uint64 { - return uint64(a.Count) + if len(selectedWireId) == 0 { + return nil, nil + } + return wireHostMap[selectedWireId][masterIdx], wireHostMap[selectedWireId][backupIdx] } diff --git a/pkg/scheduler/handler/forecast_helper.go b/pkg/scheduler/handler/forecast_helper.go index a58c630313..999f4d46db 100644 --- a/pkg/scheduler/handler/forecast_helper.go +++ b/pkg/scheduler/handler/forecast_helper.go @@ -19,7 +19,6 @@ import ( "yunion.io/x/log" - schedapi "yunion.io/x/onecloud/pkg/apis/scheduler" "yunion.io/x/onecloud/pkg/scheduler/api" "yunion.io/x/onecloud/pkg/scheduler/core" ) @@ -72,7 +71,7 @@ func transToSchedForecastResult(result *core.SchedResultItemList) interface{} { } } - items := make([]*core.SchedResultItem, 0) + items := make(core.SchedResultItems, 0) for _, item := range result.Data { hostType := item.Candidater.Getter().HostType() if schedData.Hypervisor == hostType { @@ -84,15 +83,11 @@ func transToSchedForecastResult(result *core.SchedResultItemList) interface{} { addInfos(result.Unit.LogManager.FailedLogs(), item) } - var output *schedapi.ScheduleOutput - sid := schedData.SessionId - if schedData.Backup { - output = transToBackupSchedResult(result, schedData.PreferHost, schedData.PreferBackupHost, int64(schedData.Count), sid) - } else { - output = transToRegionSchedResult(result.Data, int64(schedData.Count), sid) - } + var ( + output = transToSchedResult(result, schedData) + readyCount int64 + ) - var readyCount int64 for _, candi := range output.Candidates { if len(candi.Error) != 0 { info, exist := getOrNewFilter("select_candidate") diff --git a/pkg/scheduler/handler/handler.go b/pkg/scheduler/handler/handler.go index 29ef5ace3a..4a09e28035 100644 --- a/pkg/scheduler/handler/handler.go +++ b/pkg/scheduler/handler/handler.go @@ -124,7 +124,7 @@ func doSchedulerTest(c *gin.Context) { func transToSchedTestResult(result *core.SchedResultItemList, limit int64) interface{} { return &api.SchedTestResult{ Data: result.Data, - Total: int64(result.Len()), + Total: int64(result.Data.Len()), Limit: limit, Offset: 0, } @@ -283,15 +283,7 @@ func doSyncSchedule(c *gin.Context) { return } - count := int64(schedInfo.Count) - var resp *schedapi.ScheduleOutput - sid := schedInfo.SessionId - if schedInfo.Backup { - resp = transToBackupSchedResult(result, schedInfo.PreferHost, schedInfo.PreferBackupHost, count, sid) - } else { - resp = transToRegionSchedResult(result.Data, count, sid) - } - + resp := transToSchedResult(result, schedInfo) driver := result.Unit.GetHypervisorDriver() if err := setSchedPendingUsage(driver, schedInfo, resp); err != nil { c.AbortWithError(http.StatusInternalServerError, err) @@ -314,7 +306,16 @@ func setSchedPendingUsage(driver computemodels.IGuestDriver, req *api.SchedInfo, return nil } -func transToRegionSchedResult(result []*core.SchedResultItem, count int64, sid string) *schedapi.ScheduleOutput { +func transToSchedResult(result *core.SchedResultItemList, schedInfo *api.SchedInfo) *schedapi.ScheduleOutput { + if schedInfo.Backup { + return transToBackupSchedResult(result, + schedInfo.PreferHost, schedInfo.PreferBackupHost, int64(schedInfo.Count), schedInfo.SessionId) + } else { + return transToRegionSchedResult(result.Data, int64(schedInfo.Count), schedInfo.SessionId) + } +} + +func transToRegionSchedResult(result core.SchedResultItems, count int64, sid string) *schedapi.ScheduleOutput { apiResults := make([]*schedapi.CandidateResource, 0) succCount := 0 for _, nr := range result { diff --git a/pkg/scheduler/manager/scheduler.go b/pkg/scheduler/manager/scheduler.go index ed22453dc7..f69e3972f3 100644 --- a/pkg/scheduler/manager/scheduler.go +++ b/pkg/scheduler/manager/scheduler.go @@ -33,8 +33,8 @@ func candidatesByProvider(provider CandidatesProvider, schedData *api.SchedInfo) var err error candidateManager := provider.CandidateManager() - if len(schedData.Candidates) > 0 { - hosts, err = candidateManager.GetCandidatesByIds(provider.CandidateType(), schedData.Candidates) + if len(schedData.PreferCandidates) >= schedData.RequiredCandidates { + hosts, err = candidateManager.GetCandidatesByIds(provider.CandidateType(), schedData.PreferCandidates) } else { args := data_manager.CandidateGetArgs{ ResType: provider.CandidateType(), diff --git a/pkg/scheduler/manager/task_history.go b/pkg/scheduler/manager/task_history.go index 752e8636cd..8423e5522d 100644 --- a/pkg/scheduler/manager/task_history.go +++ b/pkg/scheduler/manager/task_history.go @@ -178,6 +178,10 @@ func (m *HistoryManager) CancelCandidatesPendingUsage(hosts []*expireHost) { continue } cancelUsage := m.GetCancelUsage(sid, hostId) + if cancelUsage == nil { + log.Errorf("failed find pending usage for session: %s, host: %s", sid, hostId) + continue + } if err := models.HostPendingUsageManager.CancelPendingUsage(hostId, cancelUsage); err != nil { log.Errorf("Cancel host %s usage %#v: %v", hostId, cancelUsage, err) } else { diff --git a/pkg/scheduler/models/pending_usage.go b/pkg/scheduler/models/pending_usage.go index bd56eae557..0aad72da81 100644 --- a/pkg/scheduler/models/pending_usage.go +++ b/pkg/scheduler/models/pending_usage.go @@ -84,6 +84,9 @@ func (m *SHostPendingUsageManager) SetPendingUsage(req *api.SchedInfo, candidate sessionUsage.StartTimer() } m.addSessionUsage(candidate.HostId, sessionUsage) + if candidate.BackupCandidate != nil { + m.SetPendingUsage(req, candidate.BackupCandidate) + } } func (m *SHostPendingUsageManager) addSessionUsage(hostId string, usage *SessionPendingUsage) {