Merge pull request #10906 from rainzm/syncprovider/localtask

fix(region): use special localtask to run time-consuming provider synchronization task
This commit is contained in:
yunion-ci-robot
2021-04-23 17:38:51 +08:00
committed by GitHub
2 changed files with 31 additions and 9 deletions
@@ -37,8 +37,8 @@ func Error2TaskData(err error) jsonutils.JSONObject {
return errJson
}
func LocalTaskRun(task ITask, proc func() (jsonutils.JSONObject, error)) {
localTaskWorkerMan.Run(func() {
func LocalTaskRunWithWorkers(task ITask, proc func() (jsonutils.JSONObject, error), wm *appsrv.SWorkerManager) {
wm.Run(func() {
log.Debugf("XXXXXXXXXXXXXXXXXXLOCAL TASK RUN STARTXXXXXXXXXXXXXXXXX")
defer log.Debugf("XXXXXXXXXXXXXXXXXXLOCAL TASK RUN END XXXXXXXXXXXXXXXXX")
@@ -59,3 +59,7 @@ func LocalTaskRun(task ITask, proc func() (jsonutils.JSONObject, error)) {
}, nil, nil)
}
func LocalTaskRun(task ITask, proc func() (jsonutils.JSONObject, error)) {
LocalTaskRunWithWorkers(task, proc, localTaskWorkerMan)
}
@@ -34,9 +34,12 @@ type CloudProviderSyncInfoTask struct {
taskman.STask
}
var syncLocalTaskWorkerMan *appsrv.SWorkerManager
func InitCloudproviderSyncWorkers(count int) {
syncWorker := appsrv.NewWorkerManager("CloudProviderSyncInfoTaskWorkerManager", count, 512, true)
taskman.RegisterTaskAndWorker(CloudProviderSyncInfoTask{}, syncWorker)
syncLocalTaskWorkerMan = appsrv.NewWorkerManager("CloudProviderSyncLocalTaskWorkerManager", count, 512, false)
}
func getAction(params *jsonutils.JSONDict) string {
@@ -94,17 +97,32 @@ func (self *CloudProviderSyncInfoTask) OnSyncCloudProviderPreInfoComplete(ctx co
syncRange := self.GetSyncRange()
db.OpsLog.LogEvent(provider, db.ACT_SYNCING_HOST, "", self.UserCred)
self.SetStage("OnSyncCloudProviderInfoComplete", nil)
provider.SyncCallSyncCloudproviderRegions(ctx, self.UserCred, syncRange)
provider.SyncCallSyncCloudproviderInterVpcNetwork(ctx, self.UserCred)
provider.CleanSchedCache()
self.SetStageComplete(ctx, nil)
db.OpsLog.LogEvent(provider, db.ACT_SYNC_HOST_COMPLETE, "", self.UserCred)
logclient.AddActionLogWithStartable(self, provider, getAction(self.Params), body, self.UserCred, true)
taskman.LocalTaskRunWithWorkers(self, func() (jsonutils.JSONObject, error) {
provider.SyncCallSyncCloudproviderRegions(ctx, self.UserCred, syncRange)
provider.SyncCallSyncCloudproviderInterVpcNetwork(ctx, self.UserCred)
return nil, nil
}, syncLocalTaskWorkerMan)
}
func (self *CloudProviderSyncInfoTask) OnSyncCloudProviderPreInfoCompleteFailed(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) {
log.Errorf("faild to sync provider quotas %s", body.String())
self.OnSyncCloudProviderPreInfoComplete(ctx, obj, body)
}
func (self *CloudProviderSyncInfoTask) OnSyncCloudProviderInfoComplete(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) {
provider := obj.(*models.SCloudprovider)
provider.CleanSchedCache()
db.OpsLog.LogEvent(provider, db.ACT_SYNC_HOST_COMPLETE, "", self.UserCred)
logclient.AddActionLogWithStartable(self, provider, getAction(self.Params), body, self.UserCred, true)
self.SetStageComplete(ctx, nil)
}
func (self *CloudProviderSyncInfoTask) OnSyncCloudProviderInfoCompleteFailed(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) {
provider := obj.(*models.SCloudprovider)
provider.CleanSchedCache()
db.OpsLog.LogEvent(provider, db.ACT_SYNC_HOST_FAILED, "", self.UserCred)
logclient.AddActionLogWithStartable(self, provider, getAction(self.Params), body, self.UserCred, false)
self.SetStageFailed(ctx, nil)
}