From 6fa0e254efa362fb35c6792035956840c56e73c7 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Sat, 16 Mar 2019 21:23:43 +0800 Subject: [PATCH 1/2] fix: cloud provider sync time is not set in auto-sync cycles --- pkg/compute/models/cloudproviders.go | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/pkg/compute/models/cloudproviders.go b/pkg/compute/models/cloudproviders.go index f82e20b401..406651598c 100644 --- a/pkg/compute/models/cloudproviders.go +++ b/pkg/compute/models/cloudproviders.go @@ -867,6 +867,9 @@ func (provider *SCloudprovider) GetEnabledCloudproviderRegions() []SCloudprovide } func (provider *SCloudprovider) syncCloudproviderRegions(userCred mcclient.TokenCredential, syncRange *SSyncRange, wg *sync.WaitGroup) { + if wg == nil { + provider.MarkSyncing(userCred) + } cprs := provider.GetEnabledCloudproviderRegions() for i := range cprs { if cprs[i].needSync() { @@ -882,6 +885,9 @@ func (provider *SCloudprovider) syncCloudproviderRegions(userCred mcclient.Token } } } + if wg == nil { + provider.MarkEndSync(userCred) + } } func (provider *SCloudprovider) SyncCallSyncCloudproviderRegions(userCred mcclient.TokenCredential, syncRange *SSyncRange) { From ca39f480580030b8f25f686a08de8a0c23d1e90d Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Mon, 18 Mar 2019 22:16:20 +0800 Subject: [PATCH 2/2] fix: more accurate sync_status update time --- pkg/compute/models/cloudaccounts.go | 32 ++++++++-- pkg/compute/models/cloudproviderregions.go | 43 +++++++++----- pkg/compute/models/cloudproviders.go | 58 ++++++++++++++----- pkg/compute/models/cloudsync.go | 2 +- pkg/compute/tasks/cloud_account_sync_task.go | 21 +++---- .../tasks/cloud_provider_sync_info_task.go | 11 ++-- 6 files changed, 116 insertions(+), 51 deletions(-) diff --git a/pkg/compute/models/cloudaccounts.go b/pkg/compute/models/cloudaccounts.go index 5f34a695f1..d3da154f1a 100644 --- a/pkg/compute/models/cloudaccounts.go +++ b/pkg/compute/models/cloudaccounts.go @@ -406,7 +406,24 @@ func (self *SCloudaccount) MarkSyncing(userCred mcclient.TokenCredential) error return nil } -func (self *SCloudaccount) MarkEndSync(userCred mcclient.TokenCredential) error { +func (self *SCloudaccount) MarkEndSyncWithLock(ctx context.Context, userCred mcclient.TokenCredential) error { + lockman.LockObject(ctx, self) + defer lockman.ReleaseObject(ctx, self) + + if self.SyncStatus == CLOUD_PROVIDER_SYNC_STATUS_IDLE { + return nil + } + + providers := self.GetCloudproviders() + for i := range providers { + if providers[i].SyncStatus != CLOUD_PROVIDER_SYNC_STATUS_IDLE { + return nil + } + } + return self.markEndSync(userCred) +} + +func (self *SCloudaccount) markEndSync(userCred mcclient.TokenCredential) error { _, err := db.Update(self, func() error { self.SyncStatus = CLOUD_PROVIDER_SYNC_STATUS_IDLE self.LastSyncEndAt = timeutils.UtcNow() @@ -1102,8 +1119,8 @@ func (account *SCloudaccount) syncAccountStatus(ctx context.Context, userCred mc account.MarkSyncing(userCred) subaccounts, err := account.probeAccountStatus(ctx, userCred) if err != nil { - account.markAccountDiscconected(ctx, userCred) account.markAllProvidersDicconnected(ctx, userCred) + account.markAccountDiscconected(ctx, userCred) return err } account.markAccountConnected(ctx, userCred) @@ -1123,16 +1140,23 @@ func (account *SCloudaccount) SubmitSyncAccountTask(ctx context.Context, userCre log.Debugf("syncAccountStatus %s %s", account.Id, account.Name) err := account.syncAccountStatus(ctx, userCred) if waitChan != nil { + if err != nil { + account.markEndSync(userCred) + } waitChan <- err } else { + syncCnt := 0 if err == nil && autoSync && account.Enabled && account.EnableAutoSync { syncRange := SSyncRange{FullSync: true} providers := account.GetEnabledCloudproviders() for i := range providers { - providers[i].syncCloudproviderRegions(userCred, &syncRange, nil) + providers[i].syncCloudproviderRegions(ctx, userCred, &syncRange, nil, autoSync) + syncCnt += 1 } } - account.MarkEndSync(userCred) + if syncCnt == 0 { + account.markEndSync(userCred) + } } }) } diff --git a/pkg/compute/models/cloudproviderregions.go b/pkg/compute/models/cloudproviderregions.go index bc3c1821c3..9bb3fbf55e 100644 --- a/pkg/compute/models/cloudproviderregions.go +++ b/pkg/compute/models/cloudproviderregions.go @@ -199,7 +199,19 @@ func (self *SCloudproviderregion) markSyncing(userCred mcclient.TokenCredential) return nil } -func (self *SCloudproviderregion) markEndSync(userCred mcclient.TokenCredential, syncResults SSyncResultSet) error { +func (self *SCloudproviderregion) markEndSync(ctx context.Context, userCred mcclient.TokenCredential, syncResults SSyncResultSet) error { + err := self.markEndSyncInternal(userCred, syncResults) + if err != nil { + return err + } + err = self.GetProvider().markEndSyncWithLock(ctx, userCred) + if err != nil { + return err + } + return nil +} + +func (self *SCloudproviderregion) markEndSyncInternal(userCred mcclient.TokenCredential, syncResults SSyncResultSet) error { _, err := db.Update(self, func() error { self.SyncStatus = CLOUD_PROVIDER_SYNC_STATUS_IDLE self.LastSyncEndAt = timeutils.UtcNow() @@ -230,11 +242,10 @@ func (set SSyncResultSet) Add(manager db.IModelManager, result compare.SyncResul } func (self *SCloudproviderregion) DoSync(ctx context.Context, userCred mcclient.TokenCredential, syncRange *SSyncRange) error { - err := self.markSyncing(userCred) - if err != nil { - log.Errorf("start sync sql fail?? %s", err) - return err - } + syncResults := SSyncResultSet{} + + self.markSyncing(userCred) + defer self.markEndSync(ctx, userCred, syncResults) localRegion := self.GetRegion() provider := self.GetProvider() @@ -244,8 +255,6 @@ func (self *SCloudproviderregion) DoSync(ctx context.Context, userCred mcclient. return err } - syncResults := SSyncResultSet{} - if localRegion.isManaged() { remoteRegion, err := driver.GetIRegionById(localRegion.ExternalId) if err == nil { @@ -255,14 +264,13 @@ func (self *SCloudproviderregion) DoSync(ctx context.Context, userCred mcclient. err = syncOnPremiseCloudProviderInfo(ctx, userCred, syncResults, provider, driver, syncRange) } + if err != nil { + log.Errorf("dosync fail %s", err) + } + log.Debugf("%s", jsonutils.Marshal(syncResults)) - err = self.markEndSync(userCred, syncResults) - if err != nil { - log.Errorf("mark end sync failed...") - return err - } - return nil + return err } func (self *SCloudproviderregion) getSyncTaskKey() string { @@ -284,7 +292,7 @@ func (self *SCloudproviderregion) submitSyncTask(userCred mcclient.TokenCredenti }) } -func (cpr *SCloudproviderregion) needSync() bool { +func (cpr *SCloudproviderregion) needAutoSync() bool { if cpr.LastSyncEndAt.IsZero() { return true } @@ -297,7 +305,10 @@ func (cpr *SCloudproviderregion) needSync() bool { isEmpty = cpr.isEmptyPublicCloud() } if isEmpty { - intval = intval * 16 // no need to check empty region + intval = intval * 16 // no need to check empty region + if intval > 24*3600 { // at least once everyday + intval = 24 * 3600 + } region := cpr.GetRegion() log.Debugf("empty region %s! no need to check so frequently", region.GetName()) } diff --git a/pkg/compute/models/cloudproviders.go b/pkg/compute/models/cloudproviders.go index 406651598c..02a02245b1 100644 --- a/pkg/compute/models/cloudproviders.go +++ b/pkg/compute/models/cloudproviders.go @@ -15,6 +15,7 @@ import ( "yunion.io/x/sqlchemy" "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/cloudprovider" "yunion.io/x/onecloud/pkg/compute/options" @@ -525,7 +526,7 @@ func (self *SCloudprovider) markStartSync(userCred mcclient.TokenCredential) err return nil } -func (self *SCloudprovider) MarkSyncing(userCred mcclient.TokenCredential) error { +func (self *SCloudprovider) markSyncing(userCred mcclient.TokenCredential) error { _, err := db.Update(self, func() error { self.SyncStatus = CLOUD_PROVIDER_SYNC_STATUS_SYNCING self.LastSync = timeutils.UtcNow() @@ -539,7 +540,34 @@ func (self *SCloudprovider) MarkSyncing(userCred mcclient.TokenCredential) error return nil } -func (self *SCloudprovider) MarkEndSync(userCred mcclient.TokenCredential) error { +func (self *SCloudprovider) markEndSyncWithLock(ctx context.Context, userCred mcclient.TokenCredential) error { + err := func() error { + lockman.LockObject(ctx, self) + defer lockman.ReleaseObject(ctx, self) + + cprs := self.GetCloudproviderRegions() + for i := range cprs { + if cprs[i].SyncStatus != CLOUD_PROVIDER_SYNC_STATUS_IDLE { + return nil + } + } + + err := self.markEndSync(userCred) + if err != nil { + return err + } + return nil + }() + + if err != nil { + return err + } + + account := self.GetCloudaccount() + return account.MarkEndSyncWithLock(ctx, userCred) +} + +func (self *SCloudprovider) markEndSync(userCred mcclient.TokenCredential) error { _, err := db.Update(self, func() error { self.SyncStatus = CLOUD_PROVIDER_SYNC_STATUS_IDLE self.LastSyncEndAt = timeutils.UtcNow() @@ -857,22 +885,22 @@ func (provider *SCloudprovider) prepareCloudproviderRegions(ctx context.Context, return cprs, nil } -func (provider *SCloudprovider) GetEnabledCloudproviderRegions() []SCloudproviderregion { +func (provider *SCloudprovider) GetCloudproviderRegions() []SCloudproviderregion { q := CloudproviderRegionManager.Query() - q = q.IsTrue("enabled") q = q.Equals("cloudprovider_id", provider.Id) - q = q.Equals("sync_status", CLOUD_PROVIDER_SYNC_STATUS_IDLE) + // q = q.IsTrue("enabled") + // q = q.Equals("sync_status", CLOUD_PROVIDER_SYNC_STATUS_IDLE) return CloudproviderRegionManager.fetchRecordsByQuery(q) } -func (provider *SCloudprovider) syncCloudproviderRegions(userCred mcclient.TokenCredential, syncRange *SSyncRange, wg *sync.WaitGroup) { - if wg == nil { - provider.MarkSyncing(userCred) - } - cprs := provider.GetEnabledCloudproviderRegions() +func (provider *SCloudprovider) syncCloudproviderRegions(ctx context.Context, userCred mcclient.TokenCredential, syncRange *SSyncRange, wg *sync.WaitGroup, autoSync bool) { + provider.markSyncing(userCred) + cprs := provider.GetCloudproviderRegions() + syncCnt := 0 for i := range cprs { - if cprs[i].needSync() { + if cprs[i].Enabled && cprs[i].CanSync() && (!autoSync || cprs[i].needAutoSync()) { + syncCnt += 1 var waitChan chan bool = nil if wg != nil { wg.Add(1) @@ -885,14 +913,14 @@ func (provider *SCloudprovider) syncCloudproviderRegions(userCred mcclient.Token } } } - if wg == nil { - provider.MarkEndSync(userCred) + if syncCnt == 0 { + provider.markEndSyncWithLock(ctx, userCred) } } -func (provider *SCloudprovider) SyncCallSyncCloudproviderRegions(userCred mcclient.TokenCredential, syncRange *SSyncRange) { +func (provider *SCloudprovider) SyncCallSyncCloudproviderRegions(ctx context.Context, userCred mcclient.TokenCredential, syncRange *SSyncRange) { var wg sync.WaitGroup - provider.syncCloudproviderRegions(userCred, syncRange, &wg) + provider.syncCloudproviderRegions(ctx, userCred, syncRange, &wg, false) wg.Wait() } diff --git a/pkg/compute/models/cloudsync.go b/pkg/compute/models/cloudsync.go index 1f21b419a9..52328310d5 100644 --- a/pkg/compute/models/cloudsync.go +++ b/pkg/compute/models/cloudsync.go @@ -24,7 +24,7 @@ type SSyncableBaseResource struct { func (self *SSyncableBaseResource) CanSync() bool { if self.SyncStatus == CLOUD_PROVIDER_SYNC_STATUS_QUEUED || self.SyncStatus == CLOUD_PROVIDER_SYNC_STATUS_SYNCING { - if self.LastSync.IsZero() || time.Now().Sub(self.LastSync) > 900*time.Second { + if self.LastSync.IsZero() || time.Now().Sub(self.LastSync) > 1800*time.Second { return true } else { return false diff --git a/pkg/compute/tasks/cloud_account_sync_task.go b/pkg/compute/tasks/cloud_account_sync_task.go index 14ad09adf9..3d936e35db 100644 --- a/pkg/compute/tasks/cloud_account_sync_task.go +++ b/pkg/compute/tasks/cloud_account_sync_task.go @@ -29,7 +29,7 @@ func (self *CloudAccountSyncInfoTask) OnInit(ctx context.Context, obj db.IStanda err := cloudaccount.SyncCallSyncAccountTask(ctx, self.UserCred) if err != nil { - cloudaccount.MarkEndSync(self.UserCred) + cloudaccount.MarkEndSyncWithLock(ctx, self.UserCred) db.OpsLog.LogEvent(cloudaccount, db.ACT_SYNC_HOST_FAILED, err, self.UserCred) self.SetStageFailed(ctx, err.Error()) logclient.AddActionLogWithStartable(self, cloudaccount, logclient.ACT_CLOUD_SYNC, err, self.UserCred, false) @@ -45,24 +45,25 @@ func (self *CloudAccountSyncInfoTask) OnInit(ctx context.Context, obj db.IStanda } if !syncRange.NeedSyncInfo() { - cloudaccount.MarkEndSync(self.UserCred) - 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) + self.OnCloudaccountSyncComplete(ctx, obj, nil) return } - self.SetStage("on_cloudaccount_sync_complete", nil) - cloudproviders := cloudaccount.GetEnabledCloudproviders() - for i := range cloudproviders { - cloudproviders[i].StartSyncCloudProviderInfoTask(ctx, self.UserCred, &syncRange, self.GetId()) + + if len(cloudproviders) > 0 { + self.SetStage("on_cloudaccount_sync_complete", nil) + for i := range cloudproviders { + cloudproviders[i].StartSyncCloudProviderInfoTask(ctx, self.UserCred, &syncRange, self.GetId()) + } + } else { + self.OnCloudaccountSyncComplete(ctx, obj, nil) } } func (self *CloudAccountSyncInfoTask) OnCloudaccountSyncComplete(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) { cloudaccount := obj.(*models.SCloudaccount) - cloudaccount.MarkEndSync(self.UserCred) + cloudaccount.MarkEndSyncWithLock(ctx, self.UserCred) 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) diff --git a/pkg/compute/tasks/cloud_provider_sync_info_task.go b/pkg/compute/tasks/cloud_provider_sync_info_task.go index e3220a41d7..ac79c53b15 100644 --- a/pkg/compute/tasks/cloud_provider_sync_info_task.go +++ b/pkg/compute/tasks/cloud_provider_sync_info_task.go @@ -50,8 +50,6 @@ func (self *CloudProviderSyncInfoTask) OnInit(ctx context.Context, obj db.IStand db.OpsLog.LogEvent(provider, db.ACT_SYNCING_HOST, "", self.UserCred) - provider.MarkSyncing(self.UserCred) - syncRange := models.SSyncRange{} syncRangeJson, _ := self.Params.Get("sync_range") if syncRangeJson != nil { @@ -59,14 +57,17 @@ func (self *CloudProviderSyncInfoTask) OnInit(ctx context.Context, obj db.IStand } syncRange.Normalize() - provider.SyncCallSyncCloudproviderRegions(self.UserCred, &syncRange) + self.SetStage("OnSyncCloudProviderInfoComplete", nil) - self.OnSyncCloudProviderInfoComplete(ctx, provider, nil) + taskman.LocalTaskRun(self, func() (jsonutils.JSONObject, error) { + provider.SyncCallSyncCloudproviderRegions(ctx, self.UserCred, &syncRange) + return nil, nil + }) } func (self *CloudProviderSyncInfoTask) OnSyncCloudProviderInfoComplete(ctx context.Context, obj db.IStandaloneModel, body jsonutils.JSONObject) { provider := obj.(*models.SCloudprovider) - provider.MarkEndSync(self.UserCred) + // provider.MarkEndSync(self.UserCred) provider.CleanSchedCache() self.SetStageComplete(ctx, nil) db.OpsLog.LogEvent(provider, db.ACT_SYNC_HOST_COMPLETE, "", self.UserCred)