From feccea5cfc5ab7752897ceba67b1509989130c35 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Sat, 10 Nov 2018 16:44:11 +0800 Subject: [PATCH] =?UTF-8?q?region:=20=E9=87=8D=E6=9E=84=E5=88=9B=E5=BB=BAs?= =?UTF-8?q?erver=E8=B0=83=E5=BA=A6=E4=BB=A3=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/compute/tasks/guest_batch_create_task.go | 103 +----------- pkg/compute/tasks/schedule.go | 164 +++++++++++++++++++ 2 files changed, 170 insertions(+), 97 deletions(-) create mode 100644 pkg/compute/tasks/schedule.go diff --git a/pkg/compute/tasks/guest_batch_create_task.go b/pkg/compute/tasks/guest_batch_create_task.go index 7e18f7b3c6..dab18475c3 100644 --- a/pkg/compute/tasks/guest_batch_create_task.go +++ b/pkg/compute/tasks/guest_batch_create_task.go @@ -2,18 +2,13 @@ package tasks import ( "context" - "fmt" "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/onecloud/pkg/cloudcommon/db" - "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" "yunion.io/x/onecloud/pkg/compute/models" - "yunion.io/x/onecloud/pkg/compute/options" - "yunion.io/x/onecloud/pkg/mcclient/auth" - "yunion.io/x/onecloud/pkg/mcclient/modules" ) type GuestBatchCreateTask struct { @@ -25,105 +20,19 @@ func init() { } func (self *GuestBatchCreateTask) OnInit(ctx context.Context, objs []db.IStandaloneModel, body jsonutils.JSONObject) { - guests := make([]*models.SGuest, 0) - for _, obj := range objs { - guest := obj.(*models.SGuest) - guests = append(guests, guest) - db.OpsLog.LogEvent(guest, db.ACT_ALLOCATING, nil, self.UserCred) - guest.SetStatus(self.UserCred, models.VM_SCHEDULE, "") - } - self.startScheduleGuests(ctx, guests) + StartScheduleObjects(ctx, self, objs) } -func (self *GuestBatchCreateTask) startScheduleGuests(ctx context.Context, guests []*models.SGuest) { - // log.Infof("%s", self.Params) - schedtags := models.ApplySchedPolicies(self.Params) - - self.SetStage("on_guest_schedule_complete", schedtags) - - s := auth.GetAdminSession(options.Options.Region, "") - results, err := modules.SchedManager.DoSchedule(s, self.Params, len(guests)) - if err != nil { - self.onSchedulerRequestFail(ctx, guests, fmt.Sprintf("Scheduler fail: %s", err)) - return - } else { - self.onSchedulerResults(ctx, guests, results) - } -} - -func (self *GuestBatchCreateTask) cancelPendingUsage(ctx context.Context) { - pendingUsage := models.SQuota{} - err := self.GetPendingUsage(&pendingUsage) - if err != nil { - log.Errorf("Taks GetPendingUsage fail %s", err) - } else { - ownerProjectId, _ := self.Params.GetString("owner_tenant_id") - err := models.QuotaManager.CancelPendingUsage(ctx, self.UserCred, ownerProjectId, &pendingUsage, &pendingUsage) - if err != nil { - log.Errorf("cancelpendingusage error %s", err) - } - } -} - -func (self *GuestBatchCreateTask) onSchedulerRequestFail(ctx context.Context, guests []*models.SGuest, reason string) { - for _, guest := range guests { - self.onScheduleFail(ctx, guest, reason) - } - self.SetStageFailed(ctx, fmt.Sprintf("Schedule failed: %s", reason)) - self.cancelPendingUsage(ctx) -} - -func (self *GuestBatchCreateTask) onScheduleFail(ctx context.Context, guest *models.SGuest, msg string) { - lockman.LockObject(ctx, guest) - defer lockman.ReleaseObject(ctx, guest) - - reason := "No matching resources" - if len(msg) > 0 { - reason = fmt.Sprintf("%s: %s", reason, msg) - } +func (self *GuestBatchCreateTask) OnScheduleFailCallback(obj IScheduleModel) { + guest := obj.(*models.SGuest) if guest.DisableDelete.IsTrue() { guest.SetDisableDelete(false) } - guest.SetStatus(self.UserCred, models.VM_SCHEDULE_FAILED, reason) - db.OpsLog.LogEvent(guest, db.ACT_ALLOCATE_FAIL, reason, self.UserCred) - notifyclient.NotifySystemError(guest.Id, guest.Name, models.VM_SCHEDULE_FAILED, reason) } -func (self *GuestBatchCreateTask) onSchedulerResults(ctx context.Context, guests []*models.SGuest, results []jsonutils.JSONObject) { - succCount := 0 - for idx := 0; idx < len(guests); idx += 1 { - guest := guests[idx] - result := results[idx] - if result.Contains("candidate") { - hostId, _ := result.GetString("candidate", "id") - self.onScheduleSucc(ctx, guest, hostId) - succCount += 1 - } else if result.Contains("error") { - msg, _ := result.Get("error") - self.onScheduleFail(ctx, guest, fmt.Sprintf("%s", msg)) - } else { - msg := fmt.Sprintf("Unknown scheduler result %s", result) - self.onScheduleFail(ctx, guest, msg) - return - } - } - if succCount == 0 { - self.SetStageFailed(ctx, "Schedule failed") - } - self.cancelPendingUsage(ctx) -} - -func (self *GuestBatchCreateTask) onScheduleSucc(ctx context.Context, guest *models.SGuest, hostId string) { - lockman.LockObject(ctx, guest) - defer lockman.ReleaseObject(ctx, guest) - - self.saveScheduleResult(ctx, guest, hostId) - models.HostManager.ClearSchedDescCache(hostId) -} - -func (self *GuestBatchCreateTask) saveScheduleResult(ctx context.Context, guest *models.SGuest, hostId string) { +func (self *GuestBatchCreateTask) SaveScheduleResult(ctx context.Context, obj IScheduleModel, hostId string) { var err error - + guest := obj.(*models.SGuest) pendingUsage := models.SQuota{} err = self.GetPendingUsage(&pendingUsage) if err != nil { @@ -185,6 +94,6 @@ func (self *GuestBatchCreateTask) saveScheduleResult(ctx context.Context, guest } } -func (self *GuestBatchCreateTask) OnGuestScheduleComplete(ctx context.Context, items []db.IStandaloneModel, data *jsonutils.JSONDict) { +func (self *GuestBatchCreateTask) OnScheduleComplete(ctx context.Context, items []db.IStandaloneModel, data *jsonutils.JSONDict) { self.SetStageComplete(ctx, nil) } diff --git a/pkg/compute/tasks/schedule.go b/pkg/compute/tasks/schedule.go new file mode 100644 index 0000000000..a21e19ee5a --- /dev/null +++ b/pkg/compute/tasks/schedule.go @@ -0,0 +1,164 @@ +package tasks + +import ( + "context" + "fmt" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" + "yunion.io/x/onecloud/pkg/cloudcommon/db/quotas" + "yunion.io/x/onecloud/pkg/cloudcommon/notifyclient" + "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/compute/options" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/mcclient/modules" +) + +const ( + SCHEDULE = models.VM_SCHEDULE + SCHEDULE_FAILED = models.VM_SCHEDULE_FAILED +) + +type IScheduleModel interface { + db.IStandaloneModel + + SetStatus(userCred mcclient.TokenCredential, status string, reason string) error +} + +type IScheduleTask interface { + GetUserCred() mcclient.TokenCredential + GetParams() *jsonutils.JSONDict + GetPendingUsage(quota quotas.IQuota) error + + SetStage(stageName string, data *jsonutils.JSONDict) + SetStageFailed(ctx context.Context, reason string) + OnScheduleFailCallback(obj IScheduleModel) + OnScheduleComplete(ctx context.Context, items []db.IStandaloneModel, data *jsonutils.JSONDict) + SaveScheduleResult(ctx context.Context, obj IScheduleModel, hostId string) +} + +func StartScheduleObjects( + ctx context.Context, + task IScheduleTask, + objs []db.IStandaloneModel, +) { + schedObjs := make([]IScheduleModel, len(objs)) + for i, obj := range objs { + schedObj := obj.(IScheduleModel) + schedObjs[i] = schedObj + db.OpsLog.LogEvent(schedObj, db.ACT_ALLOCATING, nil, task.GetUserCred()) + schedObj.SetStatus(task.GetUserCred(), SCHEDULE, "") + } + doScheduleObjects(ctx, task, schedObjs) +} + +func doScheduleObjects( + ctx context.Context, + task IScheduleTask, + objs []IScheduleModel, +) { + schedtags := models.ApplySchedPolicies(task.GetParams()) + + task.SetStage("OnScheduleComplete", schedtags) + + s := auth.GetAdminSession(options.Options.Region, "") + results, err := modules.SchedManager.DoSchedule(s, task.GetParams(), len(objs)) + if err != nil { + onSchedulerRequestFail(ctx, task, objs, fmt.Sprintf("Scheduler fail: %s", err)) + return + } + onSchedulerResults(ctx, task, objs, results) +} + +func cancelPendingUsage(ctx context.Context, task IScheduleTask) { + pendingUsage := models.SQuota{} + err := task.GetPendingUsage(&pendingUsage) + if err != nil { + log.Errorf("Taks GetPendingUsage fail %s", err) + return + } + ownerProjectId, _ := task.GetParams().GetString("owner_tenant_id") + err = models.QuotaManager.CancelPendingUsage(ctx, task.GetUserCred(), ownerProjectId, &pendingUsage, &pendingUsage) + if err != nil { + log.Errorf("cancelpendingusage error %s", err) + } +} + +func onSchedulerRequestFail( + ctx context.Context, + task IScheduleTask, + objs []IScheduleModel, + reason string, +) { + for _, obj := range objs { + onScheduleFail(ctx, task, obj, reason) + } + task.SetStageFailed(ctx, fmt.Sprintf("Schedule failed: %s", reason)) + cancelPendingUsage(ctx, task) +} + +func onScheduleFail( + ctx context.Context, + task IScheduleTask, + obj IScheduleModel, + msg string, +) { + lockman.LockObject(ctx, obj) + defer lockman.ReleaseObject(ctx, obj) + + reason := "No matching resources" + if len(msg) > 0 { + reason = fmt.Sprintf("%s: %s", reason, msg) + } + + obj.SetStatus(task.GetUserCred(), SCHEDULE_FAILED, reason) + db.OpsLog.LogEvent(obj, db.ACT_ALLOCATE_FAIL, reason, task.GetUserCred()) + notifyclient.NotifySystemError(obj.GetId(), obj.GetName(), SCHEDULE_FAILED, reason) + task.OnScheduleFailCallback(obj) +} + +func onSchedulerResults( + ctx context.Context, + task IScheduleTask, + objs []IScheduleModel, + results []jsonutils.JSONObject, +) { + succCount := 0 + for idx := 0; idx < len(objs); idx += 1 { + obj := objs[idx] + result := results[idx] + if result.Contains("candidate") { + hostId, _ := result.GetString("candidate", "id") + onScheduleSucc(ctx, task, obj, hostId) + succCount += 1 + } else if result.Contains("error") { + msg, _ := result.Get("error") + onScheduleFail(ctx, task, obj, fmt.Sprintf("%s", msg)) + } else { + msg := fmt.Sprintf("Unknown scheduler result %s", result) + onScheduleFail(ctx, task, obj, msg) + return + } + } + if succCount == 0 { + task.SetStageFailed(ctx, "Schedule failed") + } + cancelPendingUsage(ctx, task) +} + +func onScheduleSucc( + ctx context.Context, + task IScheduleTask, + obj IScheduleModel, + hostId string, +) { + lockman.LockObject(ctx, obj) + defer lockman.ReleaseObject(ctx, obj) + + task.SaveScheduleResult(ctx, obj, hostId) + models.HostManager.ClearSchedDescCache(hostId) +}