diff --git a/pkg/apis/compute/cloudaccount_const.go b/pkg/apis/compute/cloudaccount_const.go index 21ca227cb1..00a85ce0d8 100644 --- a/pkg/apis/compute/cloudaccount_const.go +++ b/pkg/apis/compute/cloudaccount_const.go @@ -23,6 +23,7 @@ const ( CLOUD_PROVIDER_DELETED = "deleted" CLOUD_PROVIDER_DELETE_FAILED = "delete_failed" + CLOUD_PROVIDER_SYNC_STATUS_QUEUING = "queuing" CLOUD_PROVIDER_SYNC_STATUS_QUEUED = "queued" CLOUD_PROVIDER_SYNC_STATUS_SYNCING = "syncing" CLOUD_PROVIDER_SYNC_STATUS_IDLE = "idle" diff --git a/pkg/cloudcommon/db/sharablebase.go b/pkg/cloudcommon/db/sharablebase.go index ba01498cf1..34fedfd38e 100644 --- a/pkg/cloudcommon/db/sharablebase.go +++ b/pkg/cloudcommon/db/sharablebase.go @@ -32,7 +32,7 @@ type SSharableBaseResourceManager struct{} func (manager *SSharableBaseResourceManager) FilterByOwner(q *sqlchemy.SQuery, owner mcclient.IIdentityProvider, scope rbacutils.TRbacScope) *sqlchemy.SQuery { if owner != nil { switch scope { - case rbacutils.ScopeDomain: + case rbacutils.ScopeProject, rbacutils.ScopeDomain: if len(owner.GetProjectDomainId()) > 0 { q = q.Filter(sqlchemy.OR( sqlchemy.Equals(q.Field("domain_id"), owner.GetProjectDomainId()), diff --git a/pkg/cloudcommon/options/options.go b/pkg/cloudcommon/options/options.go index 660acf20e0..64b1c08516 100644 --- a/pkg/cloudcommon/options/options.go +++ b/pkg/cloudcommon/options/options.go @@ -60,12 +60,12 @@ type BaseOptions struct { EnableRbac bool `help:"Switch on Role-based Access Control" default:"true"` RbacDebug bool `help:"turn on rbac debug log" default:"false"` - RbacPolicySyncPeriodSeconds int `help:"policy sync interval in seconds, default 15 minutes" default:"900"` + RbacPolicySyncPeriodSeconds int `help:"policy sync interval in seconds, default 5 minutes" default:"300"` RbacPolicySyncFailedRetrySeconds int `help:"seconds to wait after a failed sync, default 30 seconds" default:"30"` IsSlaveNode bool `help:"Region service slave node"` - CalculateQuotaUsageIntervalSeconds int `help:"interval to calculate quota usages" default:"300"` + CalculateQuotaUsageIntervalSeconds int `help:"interval to calculate quota usages, default 5 minutes" default:"300"` structarg.BaseOptions } diff --git a/pkg/compute/models/cloudaccounts.go b/pkg/compute/models/cloudaccounts.go index 1ca0873f11..6e044dfad0 100644 --- a/pkg/compute/models/cloudaccounts.go +++ b/pkg/compute/models/cloudaccounts.go @@ -26,6 +26,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/tristate" "yunion.io/x/pkg/util/timeutils" "yunion.io/x/pkg/utils" @@ -154,7 +155,7 @@ func (self *SCloudaccount) ValidateDeleteCondition(ctx context.Context) error { if self.Enabled { return httperrors.NewInvalidStatusError("account is enabled") } - if self.getSyncStatus() != api.CLOUD_PROVIDER_SYNC_STATUS_IDLE { + if self.getSyncStatus2() != api.CLOUD_PROVIDER_SYNC_STATUS_IDLE { return httperrors.NewInvalidStatusError("account is not idle") } cloudproviders := self.GetCloudproviders() @@ -471,7 +472,14 @@ func (self *SCloudaccount) markStartSync(userCred mcclient.TokenCredential) erro }) if err != nil { log.Errorf("Failed to markStartSync error: %v", err) - return err + return errors.Wrap(err, "Update") + } + providers := self.GetCloudproviders() + for i := range providers { + err := providers[i].markStartingSync(userCred) + if err != nil { + return errors.Wrap(err, "providers.markStartSync") + } } return nil } @@ -498,12 +506,10 @@ func (self *SCloudaccount) MarkEndSyncWithLock(ctx context.Context, userCred mcc return nil } - providers := self.GetCloudproviders() - for i := range providers { - if providers[i].SyncStatus != api.CLOUD_PROVIDER_SYNC_STATUS_IDLE { - return nil - } + if self.getSyncStatus2() != api.CLOUD_PROVIDER_SYNC_STATUS_IDLE { + return nil } + return self.markEndSync(userCred) } @@ -781,7 +787,7 @@ func (self *SCloudaccount) getMoreDetails(extra *jsonutils.JSONDict) *jsonutils. } extra.Add(projects, "projects") extra.Set("sync_interval_seconds", jsonutils.NewInt(int64(self.getSyncIntervalSeconds()))) - extra.Set("sync_status2", jsonutils.NewString(self.getSyncStatus())) + extra.Set("sync_status2", jsonutils.NewString(self.getSyncStatus2())) extra.Set("cloud_env", jsonutils.NewString(self.getCloudEnv())) return extra } @@ -1173,10 +1179,31 @@ func (account *SCloudaccount) needSync() bool { return false } +func (manager *SCloudaccountManager) fetchRecordsByQuery(q *sqlchemy.SQuery) []SCloudaccount { + recs := make([]SCloudaccount, 0) + err := db.FetchModelObjects(manager, q, &recs) + if err != nil { + return nil + } + return recs +} + +func (manager *SCloudaccountManager) initAllRecords() { + recs := manager.fetchRecordsByQuery(manager.Query()) + for i := range recs { + db.Update(&recs[i], func() error { + recs[i].SyncStatus = api.CLOUD_PROVIDER_SYNC_STATUS_IDLE + return nil + }) + } +} + func (manager *SCloudaccountManager) AutoSyncCloudaccountTask(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { - if isStart { - // mark all the records to be init + if isStart && !options.Options.IsSlaveNode { + // mark all the records to be idle CloudproviderRegionManager.initAllRecords() + CloudproviderManager.initAllRecords() + CloudaccountManager.initAllRecords() } q := manager.Query() @@ -1361,7 +1388,7 @@ func (self *SCloudaccount) StartCloudaccountDeleteTask(ctx context.Context, user return nil } -func (self *SCloudaccount) getSyncStatus() string { +func (self *SCloudaccount) getSyncStatus2() string { cprs := CloudproviderRegionManager.Query().SubQuery() providers := CloudproviderManager.Query().SubQuery() diff --git a/pkg/compute/models/cloudproviderregions.go b/pkg/compute/models/cloudproviderregions.go index 14785c65bb..f47887ef38 100644 --- a/pkg/compute/models/cloudproviderregions.go +++ b/pkg/compute/models/cloudproviderregions.go @@ -222,6 +222,18 @@ func (manager *SCloudproviderregionManager) FetchByIdsOrCreate(providerId string return cpr } +func (self *SCloudproviderregion) markStartingSync(userCred mcclient.TokenCredential) error { + _, err := db.Update(self, func() error { + self.SyncStatus = compute.CLOUD_PROVIDER_SYNC_STATUS_QUEUING + return nil + }) + if err != nil { + log.Errorf("Failed to markStartingSync error: %v", err) + return err + } + return nil +} + func (self *SCloudproviderregion) markStartSync(userCred mcclient.TokenCredential) error { _, err := db.Update(self, func() error { self.SyncStatus = compute.CLOUD_PROVIDER_SYNC_STATUS_QUEUED diff --git a/pkg/compute/models/cloudproviders.go b/pkg/compute/models/cloudproviders.go index 47fde56206..e97c2a2454 100644 --- a/pkg/compute/models/cloudproviders.go +++ b/pkg/compute/models/cloudproviders.go @@ -21,10 +21,9 @@ import ( "sync" "time" - "github.com/pkg/errors" - "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/tristate" "yunion.io/x/pkg/util/timeutils" "yunion.io/x/pkg/utils" @@ -248,7 +247,7 @@ func (self *SCloudprovider) syncProject(ctx context.Context, userCred mcclient.T if len(self.Name) == 0 { log.Errorf("syncProject: provider name is empty???") - return errors.New("cannot syncProject for empty name") + return errors.Error("cannot syncProject for empty name") } tenant, err := db.TenantCacheManager.FetchTenantByIdOrName(ctx, self.Name) @@ -264,7 +263,7 @@ func (self *SCloudprovider) syncProject(ctx context.Context, userCred mcclient.T params.Add(jsonutils.NewString(self.Name), "name") account := self.GetCloudaccount() if account == nil { - return errors.New("no valid cloudaccount???") + return errors.Error("no valid cloudaccount???") } domainId = account.DomainId params.Add(jsonutils.NewString(domainId), "domain_id") @@ -523,6 +522,25 @@ func (self *SCloudprovider) PerformChangeProject(ctx context.Context, userCred m return nil, self.StartSyncCloudProviderInfoTask(ctx, userCred, &SSyncRange{FullSync: true, DeepSync: true}, "") } +func (self *SCloudprovider) markStartingSync(userCred mcclient.TokenCredential) error { + _, err := db.Update(self, func() error { + self.SyncStatus = api.CLOUD_PROVIDER_SYNC_STATUS_QUEUING + return nil + }) + if err != nil { + log.Errorf("Failed to markStartSync error: %v", err) + return errors.Wrap(err, "Update") + } + cprs := self.GetCloudproviderRegions() + for i := range cprs { + err := cprs[i].markStartingSync(userCred) + if err != nil { + return errors.Wrap(err, "cprs[i].markStartingSync") + } + } + return nil +} + func (self *SCloudprovider) markStartSync(userCred mcclient.TokenCredential) error { _, err := db.Update(self, func() error { self.SyncStatus = api.CLOUD_PROVIDER_SYNC_STATUS_QUEUED @@ -532,6 +550,13 @@ func (self *SCloudprovider) markStartSync(userCred mcclient.TokenCredential) err log.Errorf("Failed to markStartSync error: %v", err) return err } + cprs := self.GetCloudproviderRegions() + for i := range cprs { + err := cprs[i].markStartingSync(userCred) + if err != nil { + return errors.Wrap(err, "cprs[i].markStartingSync") + } + } return nil } @@ -554,11 +579,12 @@ func (self *SCloudprovider) markEndSyncWithLock(ctx context.Context, userCred mc lockman.LockObject(ctx, self) defer lockman.ReleaseObject(ctx, self) - cprs := self.GetCloudproviderRegions() - for i := range cprs { - if cprs[i].SyncStatus != api.CLOUD_PROVIDER_SYNC_STATUS_IDLE { - return nil - } + if self.SyncStatus == api.CLOUD_PROVIDER_SYNC_STATUS_IDLE { + return nil + } + + if self.getSyncStatus2() != api.CLOUD_PROVIDER_SYNC_STATUS_IDLE { + return nil } err := self.markEndSync(userCred) @@ -741,6 +767,7 @@ func (self *SCloudprovider) getMoreDetails(ctx context.Context, extra *jsonutils if account != nil { extra.Add(jsonutils.NewString(account.GetName()), "cloudaccount") } + extra.Set("sync_status2", jsonutils.NewString(self.getSyncStatus2())) return extra } @@ -1131,6 +1158,7 @@ func (manager *SCloudproviderManager) FetchCustomizeColumns(ctx context.Context, func (manager *SCloudproviderManager) FilterByOwner(q *sqlchemy.SQuery, owner mcclient.IIdentityProvider, scope rbacutils.TRbacScope) *sqlchemy.SQuery { if owner != nil { + // log.Debugf("SCloudproviderManager.FilterByOwner scope:%s project:%s domain:%s", scope, owner.GetProjectId(), owner.GetProjectDomainId()) switch scope { case rbacutils.ScopeProject: if len(owner.GetProjectId()) > 0 { @@ -1159,3 +1187,38 @@ func (manager *SCloudproviderManager) FilterByOwner(q *sqlchemy.SQuery, owner mc } return q } + +func (self *SCloudprovider) getSyncStatus2() string { + q := CloudproviderRegionManager.Query() + q = q.Equals("cloudprovider_id", self.Id) + q = q.NotEquals("sync_status", api.CLOUD_PROVIDER_SYNC_STATUS_IDLE) + + cnt, err := q.CountWithError() + if err != nil { + return api.CLOUD_PROVIDER_SYNC_STATUS_ERROR + } + if cnt > 0 { + return api.CLOUD_PROVIDER_SYNC_STATUS_SYNCING + } else { + return api.CLOUD_PROVIDER_SYNC_STATUS_IDLE + } +} + +func (manager *SCloudproviderManager) fetchRecordsByQuery(q *sqlchemy.SQuery) []SCloudprovider { + recs := make([]SCloudprovider, 0) + err := db.FetchModelObjects(manager, q, &recs) + if err != nil { + return nil + } + return recs +} + +func (manager *SCloudproviderManager) initAllRecords() { + recs := manager.fetchRecordsByQuery(manager.Query()) + for i := range recs { + db.Update(&recs[i], func() error { + recs[i].SyncStatus = api.CLOUD_PROVIDER_SYNC_STATUS_IDLE + return nil + }) + } +} diff --git a/pkg/compute/models/cloudsync.go b/pkg/compute/models/cloudsync.go index 2432698e3f..2ea367f2e2 100644 --- a/pkg/compute/models/cloudsync.go +++ b/pkg/compute/models/cloudsync.go @@ -24,7 +24,7 @@ import ( "yunion.io/x/pkg/utils" "yunion.io/x/sqlchemy" - "yunion.io/x/onecloud/pkg/apis/compute" + api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -41,7 +41,7 @@ type SSyncableBaseResource struct { } func (self *SSyncableBaseResource) CanSync() bool { - if self.SyncStatus == compute.CLOUD_PROVIDER_SYNC_STATUS_QUEUED || self.SyncStatus == compute.CLOUD_PROVIDER_SYNC_STATUS_SYNCING { + if self.SyncStatus == api.CLOUD_PROVIDER_SYNC_STATUS_QUEUED || self.SyncStatus == api.CLOUD_PROVIDER_SYNC_STATUS_SYNCING { if self.LastSync.IsZero() || time.Now().Sub(self.LastSync) > 1800*time.Second { return true } else { @@ -1026,7 +1026,7 @@ func (manager *SCloudproviderregionManager) initAllRecords() { recs := manager.fetchRecordsByQuery(manager.Query()) for i := range recs { db.Update(&recs[i], func() error { - recs[i].SyncStatus = compute.CLOUD_PROVIDER_SYNC_STATUS_IDLE + recs[i].SyncStatus = api.CLOUD_PROVIDER_SYNC_STATUS_IDLE return nil }) } diff --git a/pkg/compute/models/quotas.go b/pkg/compute/models/quotas.go index 29b37cc3c4..e962e99e67 100644 --- a/pkg/compute/models/quotas.go +++ b/pkg/compute/models/quotas.go @@ -19,7 +19,7 @@ import ( "fmt" "yunion.io/x/jsonutils" - "yunion.io/x/log" + // "yunion.io/x/log" "yunion.io/x/pkg/tristate" "yunion.io/x/pkg/util/sets" @@ -105,7 +105,7 @@ func (self *SQuota) FetchUsage(ctx context.Context, scope rbacutils.TRbacScope, self.Memory = guest.TotalMemSize self.Storage = diskSize self.Eip = eipUsage.Total() - log.Debugf("%d %d %d\n", net.InternalNicCount, net.InternalVirtualNicCount, lbnic) + // log.Debugf("%d %d %d\n", net.InternalNicCount, net.InternalVirtualNicCount, lbnic) self.Port = net.InternalNicCount + net.InternalVirtualNicCount + lbnic self.Eport = net.ExternalNicCount + net.ExternalVirtualNicCount self.Bw = net.InternalBandwidth