diff --git a/pkg/compute/models/cloudaccounts.go b/pkg/compute/models/cloudaccounts.go index f0a1edc245..e406b7edc1 100644 --- a/pkg/compute/models/cloudaccounts.go +++ b/pkg/compute/models/cloudaccounts.go @@ -2307,7 +2307,7 @@ func (manager *SCloudaccountManager) AutoSyncCloudaccountStatusTask(ctx context. } cloudaccountPendingSyncs[id] = struct{}{} cloudaccountPendingSyncsMutex.Unlock() - RunSyncCloudAccountTask(ctx, func() { + RunSyncCloudAccountProbeTask(ctx, func() { defer func() { cloudaccountPendingSyncsMutex.Lock() defer cloudaccountPendingSyncsMutex.Unlock() @@ -2317,7 +2317,7 @@ func (manager *SCloudaccountManager) AutoSyncCloudaccountStatusTask(ctx context. idctx := context.WithValue(ctx, "id", id) lockman.LockObject(idctx, account) defer lockman.ReleaseObject(idctx, account) - err := account.syncAccountStatus(idctx, userCred) + err := account.syncAccountStatus(idctx, userCred, false) if err != nil { log.Errorf("unable to syncAccountStatus for cloudaccount %s: %s", account.Id, err.Error()) } @@ -2475,7 +2475,7 @@ func (acnt *SCloudaccount) setSubAccountStatus() error { return err } -func (account *SCloudaccount) syncAccountStatus(ctx context.Context, userCred mcclient.TokenCredential) error { +func (account *SCloudaccount) syncAccountStatus(ctx context.Context, userCred mcclient.TokenCredential, prepareRegions bool) error { subaccounts, err := account.probeAccountStatus(ctx, userCred) if err != nil { account.markAllProvidersDisconnected(ctx, userCred) @@ -2484,11 +2484,13 @@ func (account *SCloudaccount) syncAccountStatus(ctx context.Context, userCred mc } account.markAccountConnected(ctx, userCred) providers := account.importAllSubaccounts(ctx, userCred, subaccounts) - for i := range providers { - if providers[i].GetEnabled() { - _, err := providers[i].prepareCloudproviderRegions(ctx, userCred) - if err != nil { - providers[i].SetStatus(ctx, userCred, api.CLOUD_PROVIDER_DISCONNECTED, errors.Wrapf(err, "prepareCloudproviderRegions").Error()) + if prepareRegions { + for i := range providers { + if providers[i].GetEnabled() { + _, err := providers[i].prepareCloudproviderRegions(ctx, userCred) + if err != nil { + providers[i].SetStatus(ctx, userCred, api.CLOUD_PROVIDER_DISCONNECTED, errors.Wrapf(err, "prepareCloudproviderRegions").Error()) + } } } } @@ -2513,14 +2515,14 @@ func (account *SCloudaccount) SubmitSyncAccountTask(ctx context.Context, userCre } cloudaccountPendingSyncs[account.Id] = struct{}{} - RunSyncCloudAccountTask(ctx, func() { + RunSyncCloudAccountSyncTask(ctx, func() { defer func() { cloudaccountPendingSyncsMutex.Lock() defer cloudaccountPendingSyncsMutex.Unlock() delete(cloudaccountPendingSyncs, account.Id) }() log.Debugf("syncAccountStatus %s %s", account.Id, account.Name) - err := account.syncAccountStatus(ctx, userCred) + err := account.syncAccountStatus(ctx, userCred, true) if waitChan != nil { if err != nil { err = errors.Wrap(err, "account.syncAccountStatus") diff --git a/pkg/compute/models/syncworkers.go b/pkg/compute/models/syncworkers.go index c0a6ee91ba..e71e4626e8 100644 --- a/pkg/compute/models/syncworkers.go +++ b/pkg/compute/models/syncworkers.go @@ -33,13 +33,14 @@ import ( ) var ( - syncAccountWorker *appsrv.SWorkerManager - syncWorkers []*appsrv.SWorkerManager - syncWorkerRing *hashring.HashRing - indexMap map[string]int + syncAccountProbeWorker *appsrv.SWorkerManager + syncAccountSyncWorker *appsrv.SWorkerManager + syncWorkers []*appsrv.SWorkerManager + syncWorkerRing *hashring.HashRing + indexMap map[string]int ) -func InitSyncWorkers(count int) { +func InitSyncWorkers(count int, probeWorkerCount int, syncProbeWorkerCount int) { syncWorkers = make([]*appsrv.SWorkerManager, count) syncWorkerIndexes := make([]string, count) indexMap = map[string]int{} @@ -54,9 +55,15 @@ func InitSyncWorkers(count int) { indexMap[syncWorkerIndexes[i]] = i } syncWorkerRing = hashring.New(syncWorkerIndexes) - syncAccountWorker = appsrv.NewWorkerManager( - "cloudAccountProbeWorkerManager", - 10, + syncAccountProbeWorker = appsrv.NewWorkerManager( + "cloudAccountAutoProbeWorkerManager", + probeWorkerCount, + 2048, + true, + ) + syncAccountSyncWorker = appsrv.NewWorkerManager( + "cloudAccountSyncProbeWorkerManager", + syncProbeWorkerCount, 2048, true, ) @@ -93,14 +100,14 @@ func RunSyncCloudproviderRegionTask(ctx context.Context, key string, syncFunc fu }) } -func RunSyncCloudAccountTask(ctx context.Context, probeFunc func()) { +func runSyncCloudAccountTask(ctx context.Context, worker *appsrv.SWorkerManager, taskName string, key string, probeFunc func()) { task := resSyncTask{ syncFunc: probeFunc, - key: "AccountProb", + key: key, } - syncAccountWorker.Run(&task, nil, func(err error) { + worker.Run(&task, nil, func(err error) { data := jsonutils.NewDict() - data.Add(jsonutils.NewString("SyncCloudAccountTask"), "task_name") + data.Add(jsonutils.NewString(taskName), "task_name") data.Add(jsonutils.NewString(task.key), "task_id") data.Add(jsonutils.NewString(string(debug.Stack())), "stack") data.Add(jsonutils.NewString(err.Error()), "error") @@ -108,3 +115,15 @@ func RunSyncCloudAccountTask(ctx context.Context, probeFunc func()) { yunionconf.BugReport.SendBugReport(ctx, version.GetShortString(), string(debug.Stack()), err) }) } + +// RunSyncCloudAccountProbeTask runs auto cloud account status probe in a dedicated worker pool, +// isolated from manual sync probe tasks. +func RunSyncCloudAccountProbeTask(ctx context.Context, probeFunc func()) { + runSyncCloudAccountTask(ctx, syncAccountProbeWorker, "SyncCloudAccountProbeTask", "AccountAutoProbe", probeFunc) +} + +// RunSyncCloudAccountSyncTask runs cloud account probe before resource sync in a dedicated worker pool, +// so manual sync is not blocked by auto probe tasks. +func RunSyncCloudAccountSyncTask(ctx context.Context, probeFunc func()) { + runSyncCloudAccountTask(ctx, syncAccountSyncWorker, "SyncCloudAccountSyncTask", "AccountSyncProbe", probeFunc) +} diff --git a/pkg/compute/options/options.go b/pkg/compute/options/options.go index d52744d996..7b70993376 100644 --- a/pkg/compute/options/options.go +++ b/pkg/compute/options/options.go @@ -151,11 +151,13 @@ type ComputeOptions struct { MinimalIpAddrReusedIntervalSeconds int `help:"Minimal seconds when a release IP address can be reallocate" default:"30"` - CloudSyncWorkerCount int `help:"how many current synchronization threads" default:"5"` - CloudProviderSyncWorkerCount int `help:"how many current providers synchronize their regions, practically no limit" default:"10"` - CloudAutoSyncIntervalSeconds int `help:"frequency to check auto sync tasks" default:"300"` - DefaultSyncIntervalSeconds int `help:"minimal synchronization interval, default 15 minutes" default:"900"` - MaxCloudAccountErrorCount int `help:"maximal consecutive error count allow for a cloud account" default:"5"` + CloudSyncWorkerCount int `help:"how many current synchronization threads" default:"5"` + CloudProviderSyncWorkerCount int `help:"how many current providers synchronize their regions, practically no limit" default:"10"` + CloudAccountProbeWorkerCount int `help:"how many workers for auto cloud account status probe" default:"10"` + CloudAccountSyncProbeWorkerCount int `help:"how many workers for cloud account sync probe before resource sync" default:"10"` + CloudAutoSyncIntervalSeconds int `help:"frequency to check auto sync tasks" default:"300"` + DefaultSyncIntervalSeconds int `help:"minimal synchronization interval, default 15 minutes" default:"900"` + MaxCloudAccountErrorCount int `help:"maximal consecutive error count allow for a cloud account" default:"5"` EnableSyncName bool `help:"enable name sync" default:"true"` diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index 1d7667baa6..96552f6c26 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -135,7 +135,7 @@ func StartServiceWithJobsAndApp(jobs func(cron *cronman.SCronJobManager), appCll func startMasterTasks(opts *options.ComputeOptions, dbOpts *common_options.DBOptions, jobs func(cron *cronman.SCronJobManager)) context.CancelFunc { setInfluxdbRetentionPolicy() - models.InitSyncWorkers(opts.CloudSyncWorkerCount) + models.InitSyncWorkers(opts.CloudSyncWorkerCount, opts.CloudAccountProbeWorkerCount, opts.CloudAccountSyncProbeWorkerCount) cloudaccount_tasks.InitCloudproviderSyncWorkers(opts.CloudProviderSyncWorkerCount) var ( diff --git a/pkg/compute/tasks/cloudaccount/cloud_account_sync_task.go b/pkg/compute/tasks/cloudaccount/cloud_account_sync_task.go index b1c609a9f2..a660f5f894 100644 --- a/pkg/compute/tasks/cloudaccount/cloud_account_sync_task.go +++ b/pkg/compute/tasks/cloudaccount/cloud_account_sync_task.go @@ -128,7 +128,7 @@ func (self *CloudAccountSyncInfoTask) OnCloudaccountSyncComplete(ctx context.Con cloudaccount.MarkEndSyncWithLock(ctx, self.UserCred, syncRange.NeedSyncInfo()) db.OpsLog.LogEvent(cloudaccount, db.ACT_SYNC_HOST_COMPLETE, "", self.UserCred) self.SetStageComplete(ctx, nil) - logclient.AddActionLogWithStartable(self, cloudaccount, logclient.ACT_CLOUD_SYNC, "", self.UserCred, true) + logclient.AddActionLogWithStartable(self, cloudaccount, logclient.ACT_CLOUD_SYNC, syncRange, self.UserCred, true) } func (self *CloudAccountSyncInfoTask) OnCloudaccountSyncCompleteFailed(ctx context.Context, obj db.IStandaloneModel, err jsonutils.JSONObject) {