sync server skus from offline data (#4053)

This commit is contained in:
tb365
2019-12-10 21:22:00 +08:00
committed by Jian Qiu
parent fb1725be32
commit 7fd8192061
6 changed files with 189 additions and 505 deletions
+1 -1
View File
@@ -113,7 +113,7 @@ func syncRegionSkus(ctx context.Context, userCred mcclient.TokenCredential, loca
if cnt == 0 {
// 提前同步instance type.如果同步失败可能导致vm 内存显示为0
if err = syncSkusByRegion(localRegion); err != nil {
if err = syncServerSkusByRegion(ctx, userCred, localRegion); err != nil {
msg := fmt.Sprintf("Get Skus for region %s failed %s", localRegion.GetName(), err)
log.Errorln(msg)
// 暂时不终止同步
+98
View File
@@ -248,6 +248,10 @@ func (self *SServerSkuManager) AllowListItems(ctx context.Context, userCred mccl
return true
}
func (self SServerSku) GetGlobalId() string {
return self.ExternalId
}
func (self *SServerSku) AllowGetDetails(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject) bool {
return true
}
@@ -1071,6 +1075,12 @@ func (manager *SServerSkuManager) newFromCloudSku(ctx context.Context, userCred
return nil
}
func (manager *SServerSkuManager) newPublicCloudSku(ctx context.Context, userCred mcclient.TokenCredential, extSku SServerSku) error {
extSku.Enabled = true
extSku.Status = api.SkuStatusReady
return manager.TableSpec().Insert(&extSku)
}
func (self *SServerSku) AllowPerformEnable(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) bool {
return db.IsAllowPerform(rbacutils.ScopeSystem, userCred, self, "enable")
}
@@ -1109,6 +1119,94 @@ func (self *SServerSku) PerformDisable(ctx context.Context, userCred mcclient.To
return nil, nil
}
func (self *SServerSku) syncWithCloudSku(ctx context.Context, userCred mcclient.TokenCredential, extSku SServerSku) error {
_, err := db.Update(self, func() error {
self.PrepaidStatus = extSku.PrepaidStatus
self.PostpaidStatus = extSku.PostpaidStatus
return nil
})
return err
}
func (self *SServerSku) MarkAsSoldout(ctx context.Context) error {
_, err := db.UpdateWithLock(ctx, self, func() error {
self.PrepaidStatus = api.SkuStatusSoldout
self.PostpaidStatus = api.SkuStatusSoldout
return nil
})
return errors.Wrap(err, "SServerSku.MarkAsSoldout")
}
func (manager *SServerSkuManager) FetchSkusByRegion(regionID string) ([]SServerSku, error) {
q := manager.Query()
q = q.Equals("cloudregion_id", regionID)
skus := make([]SServerSku, 0)
err := db.FetchModelObjects(manager, q, &skus)
if err != nil {
return nil, errors.Wrap(err, "SServerSkuManager.FetchSkusByRegion")
}
return skus, nil
}
func (manager *SServerSkuManager) syncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, extSkuMeta *SSkuResourcesMeta) compare.SyncResult {
lockman.LockClass(ctx, manager, db.GetLockClassKey(manager, userCred))
defer lockman.ReleaseClass(ctx, manager, db.GetLockClassKey(manager, userCred))
syncResult := compare.SyncResult{}
extSkus, err := extSkuMeta.GetServerSkus(region)
if err != nil {
syncResult.Error(err)
return syncResult
}
dbSkus, err := manager.FetchSkusByRegion(region.GetId())
if err != nil {
syncResult.Error(err)
return syncResult
}
removed := make([]SServerSku, 0)
commondb := make([]SServerSku, 0)
commonext := make([]SServerSku, 0)
added := make([]SServerSku, 0)
err = compare.CompareSets(dbSkus, extSkus, &removed, &commondb, &commonext, &added)
if err != nil {
syncResult.Error(err)
return syncResult
}
for i := 0; i < len(removed); i += 1 {
err = removed[i].MarkAsSoldout(ctx)
if err != nil {
syncResult.DeleteError(err)
} else {
syncResult.Delete()
}
}
for i := 0; i < len(commondb); i += 1 {
err = commondb[i].syncWithCloudSku(ctx, userCred, commonext[i])
if err != nil {
syncResult.UpdateError(err)
} else {
syncResult.Update()
}
}
for i := 0; i < len(added); i += 1 {
err = manager.newPublicCloudSku(ctx, userCred, added[i])
if err != nil {
syncResult.AddError(err)
} else {
syncResult.Add()
}
}
return syncResult
}
// sku标记为soldout状态。
func (manager *SServerSkuManager) MarkAsSoldout(id string) error {
if len(id) == 0 {
-48
View File
@@ -1,48 +0,0 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package models
import (
"reflect"
"testing"
)
func Test_diff(t *testing.T) {
type args struct {
origins []string
compares []string
}
tests := []struct {
name string
args args
want []string
}{
{
name: "Test array diff",
args: args{
origins: []string{"1", "2", "3"},
compares: []string{"2", "3", "5"},
},
want: []string{"1"},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := diff(tt.args.origins, tt.args.compares); !reflect.DeepEqual(got, tt.want) {
t.Errorf("diff() = %v, want %v", got, tt.want)
}
})
}
}
+89 -11
View File
@@ -50,8 +50,9 @@ type SSkuResourcesMeta struct {
DBInstance string `json:"dbinstance"`
}
// todo: 待测试
func (self *SSkuResourcesMeta) GetServerSkus() ([]SServerSku, error) {
func (self *SSkuResourcesMeta) GetServerSkus(region *SCloudregion) ([]SServerSku, error) {
self.SetRegionFilter(region)
result := []SServerSku{}
objs, err := self.get(self.Server)
if err != nil {
@@ -63,6 +64,31 @@ func (self *SSkuResourcesMeta) GetServerSkus() ([]SServerSku, error) {
if err != nil {
return nil, errors.Wrap(err, "obj.Unmarshal")
}
// provider must not be empty
if len(sku.Provider) == 0 {
log.Debugf("source sku error: provider should not be empty. %#v", sku)
continue
}
// 处理数据
sku.Id = ""
r, err := self.fetchRegion(sku.CloudregionId)
if err != nil {
return nil, errors.Wrap(err, "SkuResourcesMeta.GetServerSkus.fetchRegion")
}
sku.CloudregionId = r.GetId()
if len(sku.ZoneId) > 0 {
zone, err := self.fetchZone(sku.ZoneId)
if err != nil {
return nil, errors.Wrap(err, "SkuResourcesMeta.GetServerSkus.fetchZone")
}
sku.ZoneId = zone.GetId()
}
result = append(result, sku)
}
return result, nil
@@ -292,21 +318,13 @@ func SyncElasticCacheSkus(ctx context.Context, userCred mcclient.TokenCredential
}
}
cloudregions := []SCloudregion{}
q := CloudregionManager.Query()
q = q.In("provider", CloudproviderManager.GetPublicProviderProvidersQuery())
err := db.FetchModelObjects(CloudregionManager, q, &cloudregions)
if err != nil {
log.Errorf("SyncElasticCacheSkus.FetchCloudregions failed: %v", err)
return
}
meta, err := fetchSkuResourcesMeta()
if err != nil {
log.Errorf("SyncElasticCacheSkus.fetchSkuResourcesMeta %s", err)
return
}
cloudregions := fetchSkuSyncCloudregions()
for i := range cloudregions {
region := &cloudregions[i]
meta.SetRegionFilter(region)
@@ -330,6 +348,53 @@ func syncElasticCacheSkusByRegion(ctx context.Context, userCred mcclient.TokenCr
log.Infof(notes)
}
// 全量同步sku列表.
func SyncServerSkus(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {
if isStart {
cnt, err := ServerSkuManager.GetPublicCloudSkuCount()
if err != nil {
log.Errorf("GetPublicCloudSkuCount fail %s", err)
return
}
if cnt > 0 {
log.Debugf("GetPublicCloudSkuCount synced skus, skip...")
return
}
}
meta, err := fetchSkuResourcesMeta()
if err != nil {
log.Errorf("SyncServerSkus.fetchSkuResourcesMeta %s", err)
return
}
cloudregions := fetchSkuSyncCloudregions()
for i := range cloudregions {
region := &cloudregions[i]
meta.SetRegionFilter(region)
result := ServerSkuManager.syncServerSkus(ctx, userCred, region, meta)
notes := fmt.Sprintf("SyncServerSkusByRegion %s result: %s", region.Name, result.Result())
log.Infof(notes)
}
// 清理无效的sku
log.Debugf("DeleteInvalidSkus in processing...")
ServerSkuManager.PendingDeleteInvalidSku()
}
// 同步指定region sku列表
func syncServerSkusByRegion(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion) error {
meta, err := fetchSkuResourcesMeta()
if err != nil {
return errors.Wrap(err, "syncServerSkusByRegion.fetchSkuResourcesMeta")
}
result := ServerSkuManager.syncServerSkus(ctx, userCred, region, meta)
notes := fmt.Sprintf("syncServerSkusByRegion %s result: %s", region.Name, result.Result())
log.Infof(notes)
return nil
}
func fetchSkuResourcesMeta() (*SSkuResourcesMeta, error) {
s := auth.GetAdminSession(context.Background(), options.Options.Region, "")
meta, err := modules.OfflineCloudmeta.GetSkuSourcesMeta(s)
@@ -345,3 +410,16 @@ func fetchSkuResourcesMeta() (*SSkuResourcesMeta, error) {
return ret, nil
}
func fetchSkuSyncCloudregions() []SCloudregion {
cloudregions := []SCloudregion{}
q := CloudregionManager.Query()
q = q.In("provider", CloudproviderManager.GetPublicProviderProvidersQuery())
err := db.FetchModelObjects(CloudregionManager, q, &cloudregions)
if err != nil {
log.Errorf("fetchSkuSyncCloudregions.FetchCloudregions failed: %v", err)
return nil
}
return cloudregions
}
-444
View File
@@ -1,444 +0,0 @@
// Copyright 2019 Yunion
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.
package models
import (
"context"
"database/sql"
"fmt"
"strings"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/compute/options"
"yunion.io/x/onecloud/pkg/mcclient"
"yunion.io/x/onecloud/pkg/mcclient/auth"
"yunion.io/x/onecloud/pkg/mcclient/modules"
)
type ServerSkus struct {
zone *SkusZone
skus []jsonutils.JSONObject
total int
updated int
created int
}
type SkusZone struct {
Provider string
RegionId string
ZoneId string
ExternalZoneId string
ExternalRegionId string
serverSkus ServerSkus
}
func mergeSkuData(odata, ndata jsonutils.JSONObject) jsonutils.JSONObject {
data, ok := odata.(*jsonutils.JSONDict)
if !ok {
log.Debugf("invalid sku dict data: %s", odata)
}
new := processSkuData(ndata)
if !ok {
log.Debugf("invalid sku dict data: %s", ndata)
}
// merge os_name
o_osname, _ := data.GetString("os_name")
if o_osname != "Any" {
n_osname, nerr := new.GetString("os_name")
if nerr != nil || n_osname != o_osname {
data.Set("os_name", jsonutils.NewString("Any"))
}
}
// merge data_disk_typs
o_disks, oerr := data.GetString("data_disk_types")
n_disks, nerr := new.GetString("data_disk_types")
if oerr != nil || nerr != nil || o_disks == "" {
data.Set("data_disk_types", jsonutils.NewString(""))
} else {
if n_disks == "" {
data.Set("data_disk_types", jsonutils.NewString(""))
} else {
data.Set("data_disk_types", jsonutils.NewString(fmt.Sprintf("%s,%s", o_disks, n_disks)))
}
}
return data
}
func processSkuData(ndata jsonutils.JSONObject) jsonutils.JSONObject {
// 从返回结果中。将os_name统一成windows|Linux|any
// 将external_id 统一替换成 id
data, ok := ndata.(*jsonutils.JSONDict)
if !ok {
log.Debugf("invalid sku dict data: %s", ndata)
}
// 处理os name
os_name, _ := ndata.GetString("os_name")
os_name = strings.ToLower(os_name)
if strings.Contains(os_name, "any") || strings.Contains(os_name, "na") || os_name == "" {
data.Set("os_name", jsonutils.NewString("Any"))
} else if os_name != "windows" {
data.Set("os_name", jsonutils.NewString("Linux"))
} else {
data.Set("os_name", jsonutils.NewString("Windows"))
}
// 将external_id 统一替换成 id.
id, err := ndata.GetString("id")
if err != nil {
data.Set("external_id", jsonutils.NewString(""))
} else {
data.Set("external_id", jsonutils.NewString(id))
data.Remove("id")
}
return data
}
func (self *ServerSkus) Init() error {
s := auth.GetAdminSession(context.Background(), options.Options.Region, "")
p, r, z := self.zone.getExternalZone()
limit := 1024
offset := 0
total := 1024
records := map[string]jsonutils.JSONObject{}
for offset < total {
ret, e := modules.CloudmetaSkus.GetSkus(s, p, r, z, limit, offset)
if e != nil {
log.Debugf("SkusZone %s init failed, %s", z, e.Error())
return e
}
for _, sku := range ret.Data {
name, err := sku.GetString("name")
if err != nil {
log.Debugf("SkusZone sku name empty : %s", sku)
return err
}
if odata, exists := records[name]; exists {
records[name] = mergeSkuData(odata, sku)
} else {
records[name] = processSkuData(sku)
}
}
offset += limit
total = ret.Total
}
filtedData := []jsonutils.JSONObject{}
for _, item := range records {
filtedData = append(filtedData, item)
}
self.total = len(records)
self.skus = filtedData
return nil
}
func (self *ServerSkus) SyncToLocalDB() error {
log.Debugf("SkusZone %s start sync.", self.zone.ExternalZoneId)
// 更新已经soldout的sku
localIds, err := ServerSkuManager.FetchAllAvailableSkuIdByZoneId(self.zone.ZoneId)
if err != nil {
return err
}
// 本次已被更新的sku id
updatedIds := make([]string, 0)
for _, sku := range self.skus {
name, _ := sku.GetString("name")
if obj, err := ServerSkuManager.FetchByZoneId(self.zone.ZoneId, name); err != nil {
if err != sql.ErrNoRows {
log.Debugf("SyncToLocalDB zone %s name %s : %s", self.zone.ZoneId, name, err.Error())
return err
}
data := SServerSku{}
if e := sku.Unmarshal(&data); e != nil {
log.Debugf("sku Unmarshal failed: %s, %s", sku, e.Error())
return e
}
if err := self.doCreate(data); err != nil {
return err
}
} else {
odata, ok := obj.(*SServerSku)
if !ok {
return fmt.Errorf("SkusZone model assertion error. %s", obj)
}
if err := self.doUpdate(odata, sku); err != nil {
return err
}
updatedIds = append(updatedIds, odata.Id)
}
}
// 处理已经下架的sku: 将本次未更新且处于available状态的sku置为soldout状态
abandonIds := diff(localIds, updatedIds)
log.Debugf("SyncToLocalDB abandon sku %s", abandonIds)
err = ServerSkuManager.MarkAllAsSoldout(abandonIds)
if err != nil {
return err
}
defer log.Debugf("SkusZone %s sync to local db.total %d,created %d,updated %d. abandoned %d", self.zone.ExternalZoneId, self.total, self.created, self.updated, len(abandonIds))
return nil
}
func (self *ServerSkus) doCreate(data SServerSku) error {
data.CloudregionId = self.zone.RegionId
data.ZoneId = self.zone.ZoneId
data.Provider = self.zone.Provider
data.Status = api.SkuStatusReady
data.Enabled = true
if err := ServerSkuManager.TableSpec().Insert(&data); err != nil {
log.Debugf("SkusZone doCreate fail: %s", err.Error())
return err
}
self.created += 1
return nil
}
func (self *ServerSkus) doUpdate(odata *SServerSku, sku jsonutils.JSONObject) error {
_, err := db.Update(odata, func() error {
if err := sku.Unmarshal(&odata); err != nil {
return err
}
odata.CloudregionId = self.zone.RegionId
odata.ZoneId = self.zone.ZoneId
odata.Provider = self.zone.Provider
// 公有云默认都是ready并启用
odata.Status = api.SkuStatusReady
odata.Enabled = true
return nil
})
if err != nil {
log.Debugf("SkusZone doUpdate fail: %s", err.Error())
return err
}
self.updated += 1
return nil
}
func (self *SkusZone) Init() error {
self.serverSkus = ServerSkus{zone: self}
err := self.serverSkus.Init()
if err != nil {
return errors.Wrap(err, "SkusZone.Init.serverSkus")
}
return nil
}
func (self *SkusZone) SyncToLocalDB() error {
err := self.serverSkus.SyncToLocalDB()
if err != nil {
return err
}
return nil
}
func (self *SkusZone) getExternalZone() (string, string, string) {
parts := strings.Split(self.ExternalZoneId, "/")
if len(parts) == 3 {
// provider, region, zone
return parts[0], parts[1], parts[2]
} else if len(parts) == 2 && parts[0] == api.CLOUD_PROVIDER_AZURE {
// azure 没有zone的概念
return parts[0], parts[1], parts[1]
}
log.Debugf("SkusZone invalid external zone id %s", self.ExternalZoneId)
return "", "", ""
}
type SkusZoneList struct {
Data []*SkusZone
total int
scuccesed int
failed int
}
func (self *SkusZoneList) initData(provider string, region SCloudregion, zones []SZone) {
for _, z := range zones {
log.Debugf("SkusZoneList initData provider %s zone %s", provider, z.GetId())
skusZone := &SkusZone{
Provider: provider,
RegionId: region.GetId(),
ZoneId: z.GetId(),
ExternalZoneId: z.GetExternalId(),
ExternalRegionId: region.GetExternalId(),
}
self.Data = append(self.Data, skusZone)
}
}
func (self *SkusZoneList) Refresh(providerIds *[]string) error {
self.Data = []*SkusZone{}
var pIds []string
if providerIds == nil {
pIds = cloudprovider.GetRegistedProviderIds()
} else {
pIds = *providerIds
}
for _, p := range pIds {
regions, e := CloudregionManager.GetRegionByProvider(p)
if e != nil {
return e
}
for _, r := range regions {
zones, e := ZoneManager.GetZonesByRegion(&r)
if e != nil {
return e
}
self.initData(p, r, zones)
}
}
self.refresh()
return nil
}
func (self *SkusZoneList) refresh() {
self.total = len(self.Data)
self.scuccesed = 0
self.failed = 0
}
func (self *SkusZoneList) SyncToLocalDB() error {
var err error
log.Debugf("######################Start Sync Skus To LocalDB######################")
for _, d := range self.Data {
if e := d.Init(); e != nil {
log.Errorf("SkusZoneList init failed: %s", e.Error())
self.failed += 1
err = e
continue
}
if e := d.SyncToLocalDB(); e != nil {
log.Errorf("SkusZoneList SyncToLocalDB failed: %s", e.Error())
self.failed += 1
err = e
continue
}
self.scuccesed += 1
remain := self.total - self.scuccesed - self.failed
log.Infof("SkusZoneList total %d, success %d.fail %d, remain %d. sync zone %s.", self.total, self.scuccesed, self.failed, remain, d.ExternalZoneId)
}
log.Debugf("######################Finished Sync Skus To LocalDB######################")
return err
}
// 全量同步sku列表.
func SyncSkus(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {
if isStart {
cnt, err := ServerSkuManager.GetPublicCloudSkuCount()
if err != nil {
log.Errorf("GetPublicCloudSkuCount fail %s", err)
return
}
if cnt > 0 {
log.Debugf("GetPublicCloudSkuCount synced skus, skip...")
return
}
}
skulist := SkusZoneList{}
if e := skulist.Refresh(nil); e != nil {
log.Errorf("SyncSkus refresh failed, %s", e.Error())
}
if e := skulist.SyncToLocalDB(); e != nil {
log.Errorf("SyncSkus sync to local db failed, %s", e.Error())
}
// 清理无效的sku
log.Debugf("DeleteInvalidSkus in processing...")
ServerSkuManager.PendingDeleteInvalidSku()
}
// 同步指定provider sku列表
func SyncSkusByProviderIds(providerIds []string) error {
skulist := SkusZoneList{}
log.Debugf("SyncSkusByProviderIds %s", providerIds)
if e := skulist.Refresh(&providerIds); e != nil {
return fmt.Errorf("SyncSkus refresh failed, %s", e.Error())
}
if e := skulist.SyncToLocalDB(); e != nil {
return fmt.Errorf("SyncSkus sync to local db failed, %s", e.Error())
}
return nil
}
// 同步指定region sku列表
func syncSkusByRegion(region *SCloudregion) error {
skulist := SkusZoneList{}
zones, err := ZoneManager.GetZonesByRegion(region)
if err != nil {
return err
}
log.Debugf("SyncSkusByRegion %s", region.GetName())
skulist.initData(region.Provider, *region, zones)
skulist.refresh()
if e := skulist.SyncToLocalDB(); e != nil {
return fmt.Errorf("SyncSkus sync to local db failed, %s", e.Error())
}
return nil
}
// 找出origins中存在,但是compares中不存在的element
func diff(origins, compares []string) []string {
ret := make([]string, 0)
for _, o := range origins {
if !utils.IsInStringArray(o, compares) && len(o) > 0 {
ret = append(ret, o)
}
}
return ret
}
+1 -1
View File
@@ -93,7 +93,7 @@ func StartService() {
cron.AddJobEveryFewHour("AutoDiskSnapshot", 1, 5, 0, models.DiskManager.AutoDiskSnapshot, false)
cron.AddJobEveryFewHour("SnapshotsCleanup", 1, 35, 0, models.SnapshotManager.CleanupSnapshots, false)
cron.AddJobEveryFewHour("AutoSyncExtDiskSnapshot", 1, 10, 0, models.DiskManager.AutoSyncExtDiskSnapshot, false)
cron.AddJobEveryFewDays("SyncSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncSkus, true)
cron.AddJobEveryFewDays("SyncSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncServerSkus, true)
cron.AddJobEveryFewDays("SyncDBInstanceSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncDBInstanceSkus, true)
cron.AddJobEveryFewDays("SyncElasticCacheSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncElasticCacheSkus, true)
cron.AddJobEveryFewDays("StorageSnapshotsRecycle", 1, 2, 0, 0, models.StorageManager.StorageSnapshotsRecycle, false)