mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-19 10:46:58 +08:00
fix(region): split sync and probe worker (#25101)
This commit is contained in:
@@ -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")
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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"`
|
||||
|
||||
|
||||
@@ -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 (
|
||||
|
||||
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user