mirror of
https://hubproxy.babadafafafafa.cn/https://github.com/yunionio/cloudpods.git
synced 2026-09-20 08:03:53 +08:00
feat(cloudcommon): db record checksum for consistency (#13679)
This commit is contained in:
@@ -102,6 +102,7 @@ func init() {
|
||||
cmd.Perform("probe-isolated-devices", &options.ServerIdOptions{})
|
||||
cmd.Perform("cpuset", &options.ServerCPUSetOptions{})
|
||||
cmd.Perform("cpuset-remove", &options.ServerIdOptions{})
|
||||
cmd.Perform("calculate-record-checksum", &options.ServerIdOptions{})
|
||||
|
||||
cmd.Get("vnc", new(options.ServerIdOptions))
|
||||
cmd.Get("desc", new(options.ServerIdOptions))
|
||||
|
||||
@@ -159,6 +159,7 @@ var (
|
||||
"sql_connection",
|
||||
"clickhouse",
|
||||
"ops_log_with_clickhouse",
|
||||
"db_checksum_tables",
|
||||
"auto_sync_table",
|
||||
"exit_after_db_init",
|
||||
"global_virtual_resource_namespace",
|
||||
|
||||
120
pkg/cloudcommon/db/checksum.go
Normal file
120
pkg/cloudcommon/db/checksum.go
Normal file
@@ -0,0 +1,120 @@
|
||||
// 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 db
|
||||
|
||||
import (
|
||||
"crypto/md5"
|
||||
"fmt"
|
||||
"reflect"
|
||||
"sort"
|
||||
|
||||
"yunion.io/x/log"
|
||||
"yunion.io/x/pkg/errors"
|
||||
"yunion.io/x/pkg/util/reflectutils"
|
||||
"yunion.io/x/pkg/utils"
|
||||
)
|
||||
|
||||
func calculateRecordChecksumByValues(vals []interface{}) string {
|
||||
ss := ""
|
||||
for _, val := range vals {
|
||||
ss += fmt.Sprintf("\n%v", val)
|
||||
}
|
||||
hStr := md5.Sum([]byte(ss))
|
||||
sum := fmt.Sprintf("%x", hStr)
|
||||
log.Debugf("calculate values string: %s checksum: %s", ss, sum)
|
||||
return sum
|
||||
}
|
||||
|
||||
func CalculateModelChecksum(dbObj IModel) (string, error) {
|
||||
if !dbObj.GetModelManager().IsEnableRecordChecksum() {
|
||||
return "", nil
|
||||
}
|
||||
|
||||
objMan := dbObj.GetModelManager()
|
||||
if objMan == nil {
|
||||
return "", errors.Errorf("Object %#v not set model manager", dbObj)
|
||||
}
|
||||
|
||||
dataValue := reflect.ValueOf(dbObj).Elem()
|
||||
dataFields := reflectutils.FetchStructFieldValueSet(dataValue)
|
||||
|
||||
cols := objMan.TableSpec().Columns()
|
||||
vals := []interface{}{}
|
||||
keys := make([]string, 0)
|
||||
for _, c := range cols {
|
||||
keys = append(keys, c.Name())
|
||||
}
|
||||
sort.Strings(keys)
|
||||
for _, k := range keys {
|
||||
|
||||
if utils.IsInStringArray(k, []string{COLUMN_RECORD_CHECKSUM, COLUMN_UPDATE_VERSION, COLUMN_UPDATED_AT}) {
|
||||
continue
|
||||
}
|
||||
v, find := dataFields.GetInterface(k)
|
||||
if !find {
|
||||
continue
|
||||
}
|
||||
vals = append(vals, v)
|
||||
}
|
||||
|
||||
return calculateRecordChecksumByValues(vals), nil
|
||||
}
|
||||
|
||||
func UpdateModelChecksum(dbObj IModel) error {
|
||||
ts := dbObj.GetModelManager().TableSpec().GetTableSpec()
|
||||
updateChecksum, err := CalculateModelChecksum(dbObj)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "CalculateModelChecksum for update")
|
||||
}
|
||||
_, err = ts.Update(dbObj, func() error {
|
||||
dbObj.SetRecordChecksum(updateChecksum)
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "UpdateModelChecksum to %s", updateChecksum)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func InjectModelsChecksum(man IModelManager) error {
|
||||
limit := 2048
|
||||
q := man.Query()
|
||||
totalCnt, err := q.CountWithError()
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "Get total records count")
|
||||
}
|
||||
setCnt := 0
|
||||
for {
|
||||
if setCnt >= totalCnt {
|
||||
break
|
||||
}
|
||||
q = q.Limit(limit).Offset(setCnt)
|
||||
objs, err := FetchIModelObjects(man, q)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "FetchModelObjects")
|
||||
}
|
||||
for i := range objs {
|
||||
obj := objs[i]
|
||||
err := UpdateModelChecksum(obj)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "UpdateModelChecksum for %s %s", man.Keyword(), obj.GetId())
|
||||
} else {
|
||||
log.Debugf("object %s %s %s calculate checksum completed", obj.Keyword(), obj.GetId(), obj.GetName())
|
||||
}
|
||||
}
|
||||
setCnt += len(objs)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -1105,6 +1105,29 @@ func FetchModelObjects(modelManager IModelManager, query *sqlchemy.SQuery, targe
|
||||
return nil
|
||||
}
|
||||
|
||||
func FetchIModelObjects(modelManager IModelManager, query *sqlchemy.SQuery) ([]IModel, error) {
|
||||
// TODO: refactor below duplicated code from FetchModelObjects
|
||||
rows, err := query.Rows()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer rows.Close()
|
||||
|
||||
objs := make([]IModel, 0)
|
||||
for rows.Next() {
|
||||
m, err := NewModelObject(modelManager)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
err = query.Row2Struct(rows, m)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
objs = append(objs, m)
|
||||
}
|
||||
return objs, nil
|
||||
}
|
||||
|
||||
func DoCreate(manager IModelManager, ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject, ownerId mcclient.IIdentityProvider) (IModel, error) {
|
||||
lockman.LockClass(ctx, manager, GetLockClassKey(manager, ownerId))
|
||||
defer lockman.ReleaseClass(ctx, manager, GetLockClassKey(manager, ownerId))
|
||||
|
||||
@@ -217,7 +217,13 @@ func fetchItem(manager IModelManager, ctx context.Context, userCred mcclient.Tok
|
||||
if err != nil {
|
||||
item, err = fetchItemByName(manager, ctx, userCred, idStr, query)
|
||||
}
|
||||
return item, err
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if err := item.CheckConsistent(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return item, nil
|
||||
}
|
||||
|
||||
func FetchUserInfo(ctx context.Context, data jsonutils.JSONObject) (mcclient.IIdentityProvider, error) {
|
||||
|
||||
@@ -43,6 +43,10 @@ type IModelManager interface {
|
||||
// Table() *sqlchemy.STable
|
||||
TableSpec() ITableSpec
|
||||
|
||||
// db record checksum
|
||||
EnableRecordChecksum() IModelManager
|
||||
IsEnableRecordChecksum() bool
|
||||
|
||||
// Keyword() string
|
||||
KeywordPlural() string
|
||||
Alias() string
|
||||
@@ -201,6 +205,10 @@ type IModel interface {
|
||||
|
||||
GetUsages() []IUsage
|
||||
GetI18N(ctx context.Context) *jsonutils.JSONDict
|
||||
|
||||
GetRecordChecksum() string
|
||||
SetRecordChecksum(checksum string)
|
||||
CheckConsistent() error
|
||||
}
|
||||
|
||||
type IResourceModelManager interface {
|
||||
|
||||
@@ -22,6 +22,7 @@ import (
|
||||
"time"
|
||||
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
"yunion.io/x/pkg/errors"
|
||||
"yunion.io/x/sqlchemy"
|
||||
|
||||
@@ -36,10 +37,18 @@ import (
|
||||
"yunion.io/x/onecloud/pkg/util/stringutils2"
|
||||
)
|
||||
|
||||
const (
|
||||
COLUMN_RECORD_CHECKSUM = "record_checksum"
|
||||
COLUMN_UPDATE_VERSION = "update_version"
|
||||
COLUMN_UPDATED_AT = "updated_at"
|
||||
)
|
||||
|
||||
type SModelBase struct {
|
||||
object.SObject
|
||||
|
||||
manager IModelManager `ignore:"true"` // pointer to modelmanager
|
||||
|
||||
RecordChecksum string `width:"256" charset:"ascii" nullable:"true" list:"user" json:"record_checksum"`
|
||||
}
|
||||
|
||||
type SModelBaseManager struct {
|
||||
@@ -50,6 +59,8 @@ type SModelBaseManager struct {
|
||||
keywordPlural string
|
||||
alias string
|
||||
aliasPlural string
|
||||
|
||||
enableRecordChecksum bool
|
||||
}
|
||||
|
||||
func NewModelBaseManager(model interface{}, tableName string, keyword string, keywordPlural string) SModelBaseManager {
|
||||
@@ -74,6 +85,16 @@ func (manager *SModelBaseManager) IsStandaloneManager() bool {
|
||||
return false
|
||||
}
|
||||
|
||||
func (manager *SModelBaseManager) EnableRecordChecksum() IModelManager {
|
||||
manager.enableRecordChecksum = true
|
||||
log.Debugf("manager %s enableRecordChecksum", manager.KeywordPlural())
|
||||
return manager.GetIModelManager()
|
||||
}
|
||||
|
||||
func (manager *SModelBaseManager) IsEnableRecordChecksum() bool {
|
||||
return manager.enableRecordChecksum
|
||||
}
|
||||
|
||||
func (manager *SModelBaseManager) GetIModelManager() IModelManager {
|
||||
virt := manager.GetVirtualObject()
|
||||
if virt == nil {
|
||||
@@ -650,3 +671,40 @@ func (model *SModelBase) GetUsages() []IUsage {
|
||||
func (model *SModelBase) GetI18N(ctx context.Context) *jsonutils.JSONDict {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (model *SModelBase) SetRecordChecksum(checksum string) {
|
||||
model.RecordChecksum = checksum
|
||||
}
|
||||
|
||||
func (model *SModelBase) GetRecordChecksum() string {
|
||||
return model.RecordChecksum
|
||||
}
|
||||
|
||||
func (model *SModelBase) CheckConsistent() error {
|
||||
obj := model.GetIModel()
|
||||
man := obj.GetModelManager()
|
||||
if !man.IsEnableRecordChecksum() {
|
||||
return nil
|
||||
}
|
||||
|
||||
calChecksum, err := CalculateModelChecksum(obj)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "CalculateModelChecksum")
|
||||
}
|
||||
savedChecksum := obj.GetRecordChecksum()
|
||||
if calChecksum != savedChecksum {
|
||||
log.Errorf("Record %s(%s) checksum changed, expected(%s) != calculated(%s)", obj.Keyword(), obj.GetId(), savedChecksum, calChecksum)
|
||||
return errors.Errorf("Record %s(%s) checksum changed, expected(%s) != calculated(%s)", obj.Keyword(), obj.GetId(), savedChecksum, calChecksum)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (model *SModelBase) PerformCalculateRecordChecksum(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
|
||||
checksum, err := CalculateModelChecksum(model.GetIModel())
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "CalculateModelChecksum")
|
||||
}
|
||||
return jsonutils.Marshal(map[string]string{
|
||||
"checksum": checksum,
|
||||
}), nil
|
||||
}
|
||||
|
||||
@@ -94,7 +94,7 @@ func mustCheckModelManager(modelMan IModelManager) {
|
||||
}
|
||||
}
|
||||
|
||||
func CheckSync(autoSync bool) bool {
|
||||
func CheckSync(autoSync bool, checksumTables []string, skipInitChecksum bool) bool {
|
||||
log.Infof("Start check database schema ...")
|
||||
inSync := true
|
||||
var err error
|
||||
@@ -143,6 +143,16 @@ func CheckSync(autoSync bool) bool {
|
||||
inSync = false
|
||||
}
|
||||
}
|
||||
|
||||
if utils.IsInStringArray(tableSpec.Name(), checksumTables) {
|
||||
modelMan.EnableRecordChecksum()
|
||||
if len(sqls) > 0 || !skipInitChecksum {
|
||||
if err := InjectModelsChecksum(modelMan); err != nil {
|
||||
log.Errorf("InjectModelsChecksum for %q error: %v", modelMan.TableSpec().Name(), err)
|
||||
return false
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return inSync
|
||||
}
|
||||
@@ -150,7 +160,7 @@ func CheckSync(autoSync bool) bool {
|
||||
func EnsureAppSyncDB(app *appsrv.Application, opt *common_options.DBOptions, modelInitDBFunc func() error) {
|
||||
// cloudcommon.InitDB(opt)
|
||||
|
||||
if !CheckSync(opt.AutoSyncTable) {
|
||||
if !CheckSync(opt.AutoSyncTable, opt.DBChecksumTables, opt.DBChecksumSkipInit) {
|
||||
log.Fatalf("database schema not in sync!")
|
||||
}
|
||||
|
||||
|
||||
@@ -49,6 +49,8 @@ type ITableSpec interface {
|
||||
Decrement(ctx context.Context, diff interface{}, target interface{}) error
|
||||
|
||||
GetSplitTable() *splitable.SSplitTableSpec
|
||||
|
||||
GetTableSpec() *sqlchemy.STableSpec
|
||||
}
|
||||
|
||||
type sTableSpec struct {
|
||||
@@ -126,25 +128,76 @@ func (ts *sTableSpec) isMarkDeleted(dt interface{}) (bool, error) {
|
||||
return obj.GetDeleted(), nil
|
||||
}
|
||||
|
||||
func (ts *sTableSpec) rejectRecordChecksumAfterInsert(obj IModel) error {
|
||||
if !obj.GetModelManager().IsEnableRecordChecksum() {
|
||||
return nil
|
||||
}
|
||||
_, err := ts.ITableSpec.Update(obj, func() error {
|
||||
checkSum, err := CalculateModelChecksum(obj)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "CalculateModelChecksum for InsertOrUpdate")
|
||||
}
|
||||
obj.SetRecordChecksum(checkSum)
|
||||
return nil
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
func (ts *sTableSpec) Insert(ctx context.Context, dt interface{}) error {
|
||||
if err := ts.ITableSpec.Insert(dt); err != nil {
|
||||
return err
|
||||
}
|
||||
ts.rejectRecordChecksumAfterInsert(dt.(IModel))
|
||||
ts.inform(ctx, dt, informer.Create)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (ts *sTableSpec) GetTableSpec() *sqlchemy.STableSpec {
|
||||
return ts.ITableSpec.(*sqlchemy.STableSpec)
|
||||
}
|
||||
|
||||
func (ts *sTableSpec) calculateRecordChecksum(dt interface{}) (string, error) {
|
||||
return "", errors.ErrNotImplemented
|
||||
}
|
||||
|
||||
func (ts *sTableSpec) InsertOrUpdate(ctx context.Context, dt interface{}) error {
|
||||
if err := ts.ITableSpec.InsertOrUpdate(dt); err != nil {
|
||||
return err
|
||||
}
|
||||
ts.rejectRecordChecksumAfterInsert(dt.(IModel))
|
||||
ts.inform(ctx, dt, informer.Create)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (ts *sTableSpec) CheckRecordChanged(dbObj IModel) error {
|
||||
return dbObj.CheckConsistent()
|
||||
}
|
||||
|
||||
func (ts *sTableSpec) Update(ctx context.Context, dt interface{}, doUpdate func() error) (sqlchemy.UpdateDiffs, error) {
|
||||
dbObj := dt.(IModel)
|
||||
isEnableRecordChecksum := dbObj.GetModelManager().IsEnableRecordChecksum()
|
||||
if isEnableRecordChecksum {
|
||||
if err := ts.CheckRecordChanged(dbObj); err != nil {
|
||||
log.Errorf("checkRecordChanged when update error: %s", err)
|
||||
return nil, errors.Wrap(err, "checkRecordChanged when update")
|
||||
}
|
||||
}
|
||||
|
||||
oldObj := jsonutils.Marshal(dt)
|
||||
diffs, err := ts.ITableSpec.Update(dt, doUpdate)
|
||||
diffs, err := ts.ITableSpec.Update(dt, func() error {
|
||||
if err := doUpdate(); err != nil {
|
||||
return err
|
||||
}
|
||||
if isEnableRecordChecksum {
|
||||
dbObj = dt.(IModel)
|
||||
updateChecksum, err := CalculateModelChecksum(dbObj)
|
||||
if err != nil {
|
||||
return errors.Wrap(err, "CalculateModelChecksum for update")
|
||||
}
|
||||
dbObj.SetRecordChecksum(updateChecksum)
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
@@ -147,7 +147,9 @@ type DBOptions struct {
|
||||
|
||||
Clickhouse string `help:"Connection string for click house"`
|
||||
|
||||
OpsLogWithClickhouse bool `help:"store operation logs with clickhouse" default:"false"`
|
||||
OpsLogWithClickhouse bool `help:"store operation logs with clickhouse" default:"false"`
|
||||
DBChecksumTables []string `help:"DB tables with record checksum for consistency"`
|
||||
DBChecksumSkipInit bool `help:"Skip DB tables with record checksum calculation when init" default:"false"`
|
||||
|
||||
AutoSyncTable bool `help:"Automatically synchronize table changes if differences are detected"`
|
||||
ExitAfterDBInit bool `help:"Exit program after db initialization" default:"false"`
|
||||
|
||||
@@ -45,7 +45,7 @@ func StartService() {
|
||||
cloudcommon.InitDB(&opts.DBOptions)
|
||||
defer cloudcommon.CloseDB()
|
||||
|
||||
if !db.CheckSync(opts.AutoSyncTable) {
|
||||
if !db.CheckSync(opts.AutoSyncTable, opts.DBChecksumTables, opts.DBChecksumSkipInit) {
|
||||
log.Fatalf("database schema not in sync!")
|
||||
}
|
||||
|
||||
|
||||
@@ -71,7 +71,7 @@ func StartService() error {
|
||||
if count == checkDBSyncRetries {
|
||||
log.Fatalf("database schema not in sync!!!")
|
||||
}
|
||||
if !db.CheckSync(false) {
|
||||
if !db.CheckSync(false, nil, true) {
|
||||
log.Errorf("database schema not in sync, wait region sync database")
|
||||
time.Sleep(2 * time.Second)
|
||||
} else {
|
||||
|
||||
@@ -47,7 +47,7 @@ func StartService() {
|
||||
InitHandlers(app)
|
||||
cloudcommon.AppDBInit(app)
|
||||
|
||||
if db.CheckSync(opts.AutoSyncTable) {
|
||||
if db.CheckSync(opts.AutoSyncTable, opts.DBChecksumTables, opts.DBChecksumSkipInit) {
|
||||
err := models.InitDB()
|
||||
if err == nil {
|
||||
app_common.ServeForeverWithCleanup(app, baseOpts, func() {
|
||||
|
||||
Reference in New Issue
Block a user