From 199a7d678e305382875f616d002791ef2d87798f Mon Sep 17 00:00:00 2001 From: Qu Xuan Date: Wed, 27 Jan 2021 14:55:58 +0800 Subject: [PATCH] fix(region): avoid worker discarded after state unchanged --- pkg/cloudcommon/db/taskman/coordinator.go | 4 ++-- pkg/cloudcommon/db/taskman/interface.go | 2 +- pkg/cloudcommon/db/taskman/tasks.go | 9 ++++++--- pkg/cloudcommon/db/taskman/worker.go | 13 ++++++++----- pkg/compute/models/guests.go | 8 +++++++- pkg/compute/models/helper.go | 8 +++----- 6 files changed, 27 insertions(+), 17 deletions(-) diff --git a/pkg/cloudcommon/db/taskman/coordinator.go b/pkg/cloudcommon/db/taskman/coordinator.go index ccdab83e71..fd4cbd3f6b 100644 --- a/pkg/cloudcommon/db/taskman/coordinator.go +++ b/pkg/cloudcommon/db/taskman/coordinator.go @@ -32,12 +32,12 @@ type BatchTaskStageFunc func(ctx context.Context, objs []db.IStandaloneModel, bo type IBatchTask interface { OnInit(ctx context.Context, objs []db.IStandaloneModel, body jsonutils.JSONObject) - ScheduleRun(data jsonutils.JSONObject) + ScheduleRun(data jsonutils.JSONObject) error } type ISingleTask interface { OnInit(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) - ScheduleRun(data jsonutils.JSONObject) + ScheduleRun(data jsonutils.JSONObject) error } var ITaskType reflect.Type diff --git a/pkg/cloudcommon/db/taskman/interface.go b/pkg/cloudcommon/db/taskman/interface.go index be89a8c804..b1e4627359 100644 --- a/pkg/cloudcommon/db/taskman/interface.go +++ b/pkg/cloudcommon/db/taskman/interface.go @@ -28,7 +28,7 @@ import ( type ITask interface { cloudcommon.IStartable - ScheduleRun(data jsonutils.JSONObject) + ScheduleRun(data jsonutils.JSONObject) error GetParams() *jsonutils.JSONDict GetUserCred() mcclient.TokenCredential GetTaskId() string diff --git a/pkg/cloudcommon/db/taskman/tasks.go b/pkg/cloudcommon/db/taskman/tasks.go index d29a44857b..b10355f9a4 100644 --- a/pkg/cloudcommon/db/taskman/tasks.go +++ b/pkg/cloudcommon/db/taskman/tasks.go @@ -121,7 +121,10 @@ func (manager *STaskManager) AllowPerformAction(ctx context.Context, userCred mc } func (manager *STaskManager) PerformAction(ctx context.Context, userCred mcclient.TokenCredential, taskId string, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) { - runTask(taskId, data) + err := runTask(taskId, data) + if err != nil { + return nil, errors.Wrapf(err, "runTask") + } resp := jsonutils.NewDict() // 'result': 'ok' resp.Add(jsonutils.NewString("ok"), "result") @@ -552,8 +555,8 @@ func execITask(taskValue reflect.Value, task *STask, odata jsonutils.JSONObject, saveRequestContextFuncValue.Call([]reflect.Value{reflect.ValueOf(&ctxData)}) } -func (task *STask) ScheduleRun(data jsonutils.JSONObject) { - runTask(task.Id, data) +func (task *STask) ScheduleRun(data jsonutils.JSONObject) error { + return runTask(task.Id, data) } func (self *STask) IsSubtask() bool { diff --git a/pkg/cloudcommon/db/taskman/worker.go b/pkg/cloudcommon/db/taskman/worker.go index 38159847ad..f791525434 100644 --- a/pkg/cloudcommon/db/taskman/worker.go +++ b/pkg/cloudcommon/db/taskman/worker.go @@ -16,9 +16,9 @@ package taskman import ( "context" + "fmt" "yunion.io/x/jsonutils" - "yunion.io/x/log" "yunion.io/x/onecloud/pkg/appsrv" "yunion.io/x/onecloud/pkg/util/panicutils" @@ -32,19 +32,22 @@ func init() { taskWorkerTable = make(map[string]*appsrv.SWorkerManager) } -func runTask(taskId string, data jsonutils.JSONObject) { +func runTask(taskId string, data jsonutils.JSONObject) error { taskName := TaskManager.getTaskName(taskId) if len(taskName) == 0 { - log.Errorf("no such task??? task_id=%s", taskId) - return + return fmt.Errorf("no such task??? task_id=%s", taskId) } worker := taskWorkMan if workerMan, ok := taskWorkerTable[taskName]; ok { worker = workerMan } - worker.Run(func() { + isOk := worker.Run(func() { TaskManager.execTask(taskId, data) }, nil, func(err error) { panicutils.SendPanicMessage(context.TODO(), err) }) + if !isOk { + return fmt.Errorf("worker %s(%s) not running may be droped", taskName, taskId) + } + return nil } diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index a2a80bcef0..8114e4f9fa 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -1888,7 +1888,13 @@ func (manager *SGuestManager) OnCreateComplete(ctx context.Context, items []db.I manager.SetPropertiesWithInstanceSnapshot(ctx, userCred, input.InstanceSnapshotId, items) } pendingUsage, pendingRegionUsage := getGuestResourceRequirements(ctx, userCred, input, ownerId, len(items), input.Backup) - RunBatchCreateTask(ctx, items, userCred, data, pendingUsage, pendingRegionUsage, "GuestBatchCreateTask", input.ParentTaskId) + err := RunBatchCreateTask(ctx, items, userCred, data, pendingUsage, pendingRegionUsage, "GuestBatchCreateTask", input.ParentTaskId) + if err != nil { + for i := range items { + guest := items[i].(*SGuest) + guest.SetStatus(userCred, api.VM_CREATE_FAILED, err.Error()) + } + } } func (guest *SGuest) GetGroups() []SGroupguest { diff --git a/pkg/compute/models/helper.go b/pkg/compute/models/helper.go index 5d05ec273b..f3142c2e7d 100644 --- a/pkg/compute/models/helper.go +++ b/pkg/compute/models/helper.go @@ -19,7 +19,6 @@ import ( "database/sql" "yunion.io/x/jsonutils" - "yunion.io/x/log" "yunion.io/x/pkg/errors" "yunion.io/x/pkg/utils" @@ -39,7 +38,7 @@ func RunBatchCreateTask( pendingRegionUsage SRegionQuota, taskName string, parentTaskId string, -) { +) error { taskItems := make([]db.IStandaloneModel, len(items)) for i, t := range items { taskItems[i] = t.(db.IStandaloneModel) @@ -47,10 +46,9 @@ func RunBatchCreateTask( params := data.(*jsonutils.JSONDict) task, err := taskman.TaskManager.NewParallelTask(ctx, taskName, taskItems, userCred, params, parentTaskId, "", &pendingUsage, &pendingRegionUsage) if err != nil { - log.Errorf("%s newTask error %s", taskName, err) - } else { - task.ScheduleRun(nil) + return errors.Wrapf(err, "NewParallelTask %s", taskName) } + return task.ScheduleRun(nil) } func ValidateScheduleCreateData(ctx context.Context, userCred mcclient.TokenCredential, input *api.ServerCreateInput, hypervisor string) (*api.ServerCreateInput, error) {