Kill gorp (#19786)
* Kill gorp Gorp is dead. Long live Gorp. ```release-note NONE ```
Этот коммит содержится в:
коммит произвёл
GitHub
родитель
c3ec0a145b
Коммит
4da98cb51a
@@ -6,17 +6,11 @@ package sqlstore
|
||||
import (
|
||||
"bytes"
|
||||
"database/sql/driver"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"strconv"
|
||||
"strings"
|
||||
|
||||
"github.com/dyatlov/go-opengraph/opengraph"
|
||||
"github.com/mattermost/gorp"
|
||||
"github.com/mattermost/mattermost-server/v6/model"
|
||||
"github.com/mattermost/mattermost-server/v6/shared/i18n"
|
||||
"github.com/mattermost/mattermost-server/v6/shared/mlog"
|
||||
"github.com/pkg/errors"
|
||||
)
|
||||
|
||||
type jsonArray []string
|
||||
@@ -68,85 +62,3 @@ func (t *TraceOnAdapter) Printf(format string, v ...interface{}) {
|
||||
type JSONSerializable interface {
|
||||
ToJSON() string
|
||||
}
|
||||
|
||||
type mattermConverter struct{}
|
||||
|
||||
func (me mattermConverter) ToDb(val interface{}) (interface{}, error) {
|
||||
switch t := val.(type) {
|
||||
case model.StringMap:
|
||||
return model.MapToJSON(t), nil
|
||||
case map[string]string:
|
||||
return model.MapToJSON(model.StringMap(t)), nil
|
||||
case model.StringArray:
|
||||
return model.ArrayToJSON(t), nil
|
||||
case model.StringInterface:
|
||||
return model.StringInterfaceToJSON(t), nil
|
||||
case map[string]interface{}:
|
||||
return model.StringInterfaceToJSON(model.StringInterface(t)), nil
|
||||
case JSONSerializable:
|
||||
return t.ToJSON(), nil
|
||||
case *opengraph.OpenGraph:
|
||||
return json.Marshal(t)
|
||||
case *model.PostImage:
|
||||
return json.Marshal(t)
|
||||
}
|
||||
|
||||
return val, nil
|
||||
}
|
||||
|
||||
func (me mattermConverter) FromDb(target interface{}) (gorp.CustomScanner, bool) {
|
||||
switch target.(type) {
|
||||
case *model.StringMap:
|
||||
binder := func(holder, target interface{}) error {
|
||||
s, ok := holder.(*string)
|
||||
if !ok {
|
||||
return errors.New(i18n.T("store.sql.convert_string_map"))
|
||||
}
|
||||
b := []byte(*s)
|
||||
return json.Unmarshal(b, target)
|
||||
}
|
||||
return gorp.CustomScanner{Holder: new(string), Target: target, Binder: binder}, true
|
||||
case *map[string]string:
|
||||
binder := func(holder, target interface{}) error {
|
||||
s, ok := holder.(*string)
|
||||
if !ok {
|
||||
return errors.New(i18n.T("store.sql.convert_string_map"))
|
||||
}
|
||||
b := []byte(*s)
|
||||
return json.Unmarshal(b, target)
|
||||
}
|
||||
return gorp.CustomScanner{Holder: new(string), Target: target, Binder: binder}, true
|
||||
case *model.StringArray:
|
||||
binder := func(holder, target interface{}) error {
|
||||
s, ok := holder.(*string)
|
||||
if !ok {
|
||||
return errors.New(i18n.T("store.sql.convert_string_array"))
|
||||
}
|
||||
b := []byte(*s)
|
||||
return json.Unmarshal(b, target)
|
||||
}
|
||||
return gorp.CustomScanner{Holder: new(string), Target: target, Binder: binder}, true
|
||||
case *model.StringInterface:
|
||||
binder := func(holder, target interface{}) error {
|
||||
s, ok := holder.(*string)
|
||||
if !ok {
|
||||
return errors.New(i18n.T("store.sql.convert_string_interface"))
|
||||
}
|
||||
b := []byte(*s)
|
||||
return json.Unmarshal(b, target)
|
||||
}
|
||||
return gorp.CustomScanner{Holder: new(string), Target: target, Binder: binder}, true
|
||||
case *map[string]interface{}:
|
||||
binder := func(holder, target interface{}) error {
|
||||
s, ok := holder.(*string)
|
||||
if !ok {
|
||||
return errors.New(i18n.T("store.sql.convert_string_interface"))
|
||||
}
|
||||
b := []byte(*s)
|
||||
return json.Unmarshal(b, target)
|
||||
}
|
||||
return gorp.CustomScanner{Holder: new(string), Target: target, Binder: binder}, true
|
||||
}
|
||||
|
||||
return gorp.CustomScanner{}, false
|
||||
}
|
||||
|
||||
@@ -5,8 +5,6 @@ package sqlstore
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/mattermost/gorp"
|
||||
)
|
||||
|
||||
// storeContextKey is the base type for all context keys for the store.
|
||||
@@ -35,14 +33,6 @@ func hasMaster(ctx context.Context) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
// DBFromContext is a helper utility that returns the DB handle from a given context.
|
||||
func (ss *SqlStore) DBFromContext(ctx context.Context) *gorp.DbMap {
|
||||
if hasMaster(ctx) {
|
||||
return ss.GetMaster()
|
||||
}
|
||||
return ss.GetReplica()
|
||||
}
|
||||
|
||||
// DBXFromContext is a helper utility that returns the sqlx DB handle from a given context.
|
||||
func (ss *SqlStore) DBXFromContext(ctx context.Context) *sqlxDBWrapper {
|
||||
if hasMaster(ctx) {
|
||||
|
||||
@@ -58,22 +58,23 @@ func TestDeleteUnusedFeatures(t *testing.T) {
|
||||
ss.Preference().(*SqlPreferenceStore).deleteUnusedFeatures()
|
||||
|
||||
//make sure features with value "false" have actually been deleted from the database
|
||||
if val, err := ss.Preference().(*SqlPreferenceStore).GetReplica().SelectInt(`SELECT COUNT(*)
|
||||
var val int64
|
||||
if err := ss.Preference().(*SqlPreferenceStore).GetReplicaX().Get(&val, `SELECT COUNT(*)
|
||||
FROM Preferences
|
||||
WHERE Category = :Category
|
||||
AND Value = :Val
|
||||
AND Name LIKE '`+store.FeatureTogglePrefix+`%'`, map[string]interface{}{"Category": model.PreferenceCategoryAdvancedSettings, "Val": "false"}); err != nil {
|
||||
WHERE Category = ?
|
||||
AND Value = ?
|
||||
AND Name LIKE '`+store.FeatureTogglePrefix+`%'`, model.PreferenceCategoryAdvancedSettings, "false"); err != nil {
|
||||
require.NoError(t, err)
|
||||
} else if val != 0 {
|
||||
require.Fail(t, "Found %d features with value 'false', expected all to be deleted", val)
|
||||
}
|
||||
//
|
||||
// make sure features with value "true" remain saved
|
||||
if val, err := ss.Preference().(*SqlPreferenceStore).GetReplica().SelectInt(`SELECT COUNT(*)
|
||||
if err := ss.Preference().(*SqlPreferenceStore).GetReplicaX().Get(&val, `SELECT COUNT(*)
|
||||
FROM Preferences
|
||||
WHERE Category = :Category
|
||||
AND Value = :Val
|
||||
AND Name LIKE '`+store.FeatureTogglePrefix+`%'`, map[string]interface{}{"Category": model.PreferenceCategoryAdvancedSettings, "Val": "true"}); err != nil {
|
||||
WHERE Category = ?
|
||||
AND Value = ?
|
||||
AND Name LIKE '`+store.FeatureTogglePrefix+`%'`, model.PreferenceCategoryAdvancedSettings, "true"); err != nil {
|
||||
require.NoError(t, err)
|
||||
} else if val == 0 {
|
||||
require.Fail(t, "Found %d features with value 'true', expected to find at least %d features", val, 2)
|
||||
|
||||
@@ -14,7 +14,6 @@ import (
|
||||
|
||||
"github.com/jmoiron/sqlx"
|
||||
|
||||
"github.com/mattermost/gorp"
|
||||
"github.com/mattermost/mattermost-server/v6/model"
|
||||
"github.com/mattermost/mattermost-server/v6/shared/mlog"
|
||||
"github.com/mattermost/mattermost-server/v6/store/storetest"
|
||||
@@ -28,10 +27,6 @@ func NewStoreTestWrapper(orig *SqlStore) *StoreTestWrapper {
|
||||
return &StoreTestWrapper{orig}
|
||||
}
|
||||
|
||||
func (w *StoreTestWrapper) GetMaster() *gorp.DbMap {
|
||||
return w.orig.GetMaster()
|
||||
}
|
||||
|
||||
func (w *StoreTestWrapper) GetMasterX() storetest.SqlXExecutor {
|
||||
return w.orig.GetMasterX()
|
||||
}
|
||||
@@ -251,6 +246,18 @@ func (w *sqlxTxWrapper) Exec(query string, args ...interface{}) (sql.Result, err
|
||||
return w.ExecRaw(query, args...)
|
||||
}
|
||||
|
||||
func (w *sqlxTxWrapper) ExecNoTimeout(query string, args ...interface{}) (sql.Result, error) {
|
||||
query = w.Tx.Rebind(query)
|
||||
|
||||
if w.trace {
|
||||
defer func(then time.Time) {
|
||||
printArgs(query, time.Since(then), args)
|
||||
}(time.Now())
|
||||
}
|
||||
|
||||
return w.Tx.ExecContext(context.Background(), query, args...)
|
||||
}
|
||||
|
||||
// ExecRaw is like Exec but without any rebinding of params. You need to pass
|
||||
// the exact param types of your target database.
|
||||
func (w *sqlxTxWrapper) ExecRaw(query string, args ...interface{}) (sql.Result, error) {
|
||||
|
||||
@@ -27,7 +27,6 @@ import (
|
||||
_ "github.com/golang-migrate/migrate/v4/source/file"
|
||||
"github.com/jmoiron/sqlx"
|
||||
"github.com/lib/pq"
|
||||
"github.com/mattermost/gorp"
|
||||
mbindata "github.com/mattermost/morph/sources/go_bindata"
|
||||
"github.com/pkg/errors"
|
||||
|
||||
@@ -113,13 +112,10 @@ type SqlStore struct {
|
||||
rrCounter int64
|
||||
srCounter int64
|
||||
|
||||
master *gorp.DbMap
|
||||
masterX *sqlxDBWrapper
|
||||
|
||||
Replicas []*gorp.DbMap
|
||||
ReplicaXs []*sqlxDBWrapper
|
||||
|
||||
searchReplicas []*gorp.DbMap
|
||||
searchReplicaXs []*sqlxDBWrapper
|
||||
|
||||
replicaLagHandles []*dbsql.DB
|
||||
@@ -248,34 +244,6 @@ func setupConnection(connType string, dataSource string, settings *model.SqlSett
|
||||
return db
|
||||
}
|
||||
|
||||
func getDBMap(settings *model.SqlSettings, db *dbsql.DB) *gorp.DbMap {
|
||||
connectionTimeout := time.Duration(*settings.QueryTimeout) * time.Second
|
||||
var dbMap *gorp.DbMap
|
||||
switch *settings.DriverName {
|
||||
case model.DatabaseDriverMysql:
|
||||
dbMap = &gorp.DbMap{
|
||||
Db: db,
|
||||
TypeConverter: mattermConverter{},
|
||||
Dialect: gorp.MySQLDialect{Engine: "InnoDB", Encoding: "UTF8MB4"},
|
||||
QueryTimeout: connectionTimeout,
|
||||
}
|
||||
case model.DatabaseDriverPostgres:
|
||||
dbMap = &gorp.DbMap{
|
||||
Db: db,
|
||||
TypeConverter: mattermConverter{},
|
||||
Dialect: gorp.PostgresDialect{},
|
||||
QueryTimeout: connectionTimeout,
|
||||
}
|
||||
default:
|
||||
mlog.Fatal("Failed to create dialect specific driver")
|
||||
return nil
|
||||
}
|
||||
if settings.Trace != nil && *settings.Trace {
|
||||
dbMap.TraceOn("sql-trace:", &TraceOnAdapter{})
|
||||
}
|
||||
return dbMap
|
||||
}
|
||||
|
||||
func (ss *SqlStore) SetContext(context context.Context) {
|
||||
ss.context = context
|
||||
}
|
||||
@@ -300,7 +268,6 @@ func (ss *SqlStore) initConnection() {
|
||||
}
|
||||
|
||||
handle := setupConnection("master", dataSource, ss.settings)
|
||||
ss.master = getDBMap(ss.settings, handle)
|
||||
ss.masterX = newSqlxDBWrapper(sqlx.NewDb(handle, ss.DriverName()),
|
||||
time.Duration(*ss.settings.QueryTimeout)*time.Second,
|
||||
*ss.settings.Trace)
|
||||
@@ -309,11 +276,9 @@ func (ss *SqlStore) initConnection() {
|
||||
}
|
||||
|
||||
if len(ss.settings.DataSourceReplicas) > 0 {
|
||||
ss.Replicas = make([]*gorp.DbMap, len(ss.settings.DataSourceReplicas))
|
||||
ss.ReplicaXs = make([]*sqlxDBWrapper, len(ss.settings.DataSourceReplicas))
|
||||
for i, replica := range ss.settings.DataSourceReplicas {
|
||||
handle := setupConnection(fmt.Sprintf("replica-%v", i), replica, ss.settings)
|
||||
ss.Replicas[i] = getDBMap(ss.settings, handle)
|
||||
ss.ReplicaXs[i] = newSqlxDBWrapper(sqlx.NewDb(handle, ss.DriverName()),
|
||||
time.Duration(*ss.settings.QueryTimeout)*time.Second,
|
||||
*ss.settings.Trace)
|
||||
@@ -324,11 +289,9 @@ func (ss *SqlStore) initConnection() {
|
||||
}
|
||||
|
||||
if len(ss.settings.DataSourceSearchReplicas) > 0 {
|
||||
ss.searchReplicas = make([]*gorp.DbMap, len(ss.settings.DataSourceSearchReplicas))
|
||||
ss.searchReplicaXs = make([]*sqlxDBWrapper, len(ss.settings.DataSourceSearchReplicas))
|
||||
for i, replica := range ss.settings.DataSourceSearchReplicas {
|
||||
handle := setupConnection(fmt.Sprintf("search-replica-%v", i), replica, ss.settings)
|
||||
ss.searchReplicas[i] = getDBMap(ss.settings, handle)
|
||||
ss.searchReplicaXs[i] = newSqlxDBWrapper(sqlx.NewDb(handle, ss.DriverName()),
|
||||
time.Duration(*ss.settings.QueryTimeout)*time.Second,
|
||||
*ss.settings.Trace)
|
||||
@@ -386,10 +349,6 @@ func (ss *SqlStore) GetDbVersion(numerical bool) (string, error) {
|
||||
|
||||
}
|
||||
|
||||
func (ss *SqlStore) GetMaster() *gorp.DbMap {
|
||||
return ss.master
|
||||
}
|
||||
|
||||
func (ss *SqlStore) GetMasterX() *sqlxDBWrapper {
|
||||
return ss.masterX
|
||||
}
|
||||
@@ -403,22 +362,6 @@ func (ss *SqlStore) SetMasterX(db *sql.DB) {
|
||||
}
|
||||
}
|
||||
|
||||
func (ss *SqlStore) GetSearchReplica() *gorp.DbMap {
|
||||
ss.licenseMutex.RLock()
|
||||
license := ss.license
|
||||
ss.licenseMutex.RUnlock()
|
||||
if license == nil {
|
||||
return ss.GetMaster()
|
||||
}
|
||||
|
||||
if len(ss.settings.DataSourceSearchReplicas) == 0 {
|
||||
return ss.GetReplica()
|
||||
}
|
||||
|
||||
rrNum := atomic.AddInt64(&ss.srCounter, 1) % int64(len(ss.searchReplicas))
|
||||
return ss.searchReplicas[rrNum]
|
||||
}
|
||||
|
||||
func (ss *SqlStore) GetSearchReplicaX() *sqlxDBWrapper {
|
||||
ss.licenseMutex.RLock()
|
||||
license := ss.license
|
||||
@@ -435,18 +378,6 @@ func (ss *SqlStore) GetSearchReplicaX() *sqlxDBWrapper {
|
||||
return ss.searchReplicaXs[rrNum]
|
||||
}
|
||||
|
||||
func (ss *SqlStore) GetReplica() *gorp.DbMap {
|
||||
ss.licenseMutex.RLock()
|
||||
license := ss.license
|
||||
ss.licenseMutex.RUnlock()
|
||||
if len(ss.settings.DataSourceReplicas) == 0 || ss.lockedToMaster || license == nil {
|
||||
return ss.GetMaster()
|
||||
}
|
||||
|
||||
rrNum := atomic.AddInt64(&ss.rrCounter, 1) % int64(len(ss.Replicas))
|
||||
return ss.Replicas[rrNum]
|
||||
}
|
||||
|
||||
func (ss *SqlStore) GetReplicaX() *sqlxDBWrapper {
|
||||
ss.licenseMutex.RLock()
|
||||
license := ss.license
|
||||
@@ -455,7 +386,7 @@ func (ss *SqlStore) GetReplicaX() *sqlxDBWrapper {
|
||||
return ss.GetMasterX()
|
||||
}
|
||||
|
||||
rrNum := atomic.AddInt64(&ss.rrCounter, 1) % int64(len(ss.Replicas))
|
||||
rrNum := atomic.AddInt64(&ss.rrCounter, 1) % int64(len(ss.ReplicaXs))
|
||||
return ss.ReplicaXs[rrNum]
|
||||
}
|
||||
|
||||
@@ -507,8 +438,8 @@ func (ss *SqlStore) TotalReadDbConnections() int {
|
||||
}
|
||||
|
||||
count := 0
|
||||
for _, db := range ss.Replicas {
|
||||
count = count + db.Db.Stats().OpenConnections
|
||||
for _, db := range ss.ReplicaXs {
|
||||
count = count + db.Stats().OpenConnections
|
||||
}
|
||||
|
||||
return count
|
||||
@@ -520,8 +451,8 @@ func (ss *SqlStore) TotalSearchDbConnections() int {
|
||||
}
|
||||
|
||||
count := 0
|
||||
for _, db := range ss.searchReplicas {
|
||||
count = count + db.Db.Stats().OpenConnections
|
||||
for _, db := range ss.searchReplicaXs {
|
||||
count = count + db.Stats().OpenConnections
|
||||
}
|
||||
|
||||
return count
|
||||
@@ -748,10 +679,10 @@ func IsUniqueConstraintError(err error, indexName []string) bool {
|
||||
return unique && field
|
||||
}
|
||||
|
||||
func (ss *SqlStore) GetAllConns() []*gorp.DbMap {
|
||||
all := make([]*gorp.DbMap, len(ss.Replicas)+1)
|
||||
copy(all, ss.Replicas)
|
||||
all[len(ss.Replicas)] = ss.master
|
||||
func (ss *SqlStore) GetAllConns() []*sqlxDBWrapper {
|
||||
all := make([]*sqlxDBWrapper, len(ss.ReplicaXs)+1)
|
||||
copy(all, ss.ReplicaXs)
|
||||
all[len(ss.ReplicaXs)] = ss.masterX
|
||||
return all
|
||||
}
|
||||
|
||||
@@ -762,24 +693,24 @@ func (ss *SqlStore) RecycleDBConnections(d time.Duration) {
|
||||
originalDuration := time.Duration(*ss.settings.ConnMaxLifetimeMilliseconds) * time.Millisecond
|
||||
// Set the max lifetimes for all connections.
|
||||
for _, conn := range ss.GetAllConns() {
|
||||
conn.Db.SetConnMaxLifetime(d)
|
||||
conn.SetConnMaxLifetime(d)
|
||||
}
|
||||
// Wait for that period with an additional 2 seconds of scheduling delay.
|
||||
time.Sleep(d + 2*time.Second)
|
||||
// Reset max lifetime back to original value.
|
||||
for _, conn := range ss.GetAllConns() {
|
||||
conn.Db.SetConnMaxLifetime(originalDuration)
|
||||
conn.SetConnMaxLifetime(originalDuration)
|
||||
}
|
||||
}
|
||||
|
||||
func (ss *SqlStore) Close() {
|
||||
ss.master.Db.Close()
|
||||
for _, replica := range ss.Replicas {
|
||||
replica.Db.Close()
|
||||
ss.masterX.Close()
|
||||
for _, replica := range ss.ReplicaXs {
|
||||
replica.Close()
|
||||
}
|
||||
|
||||
for _, replica := range ss.searchReplicas {
|
||||
replica.Db.Close()
|
||||
for _, replica := range ss.searchReplicaXs {
|
||||
replica.Close()
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -17,7 +17,6 @@ import (
|
||||
|
||||
"github.com/go-sql-driver/mysql"
|
||||
"github.com/lib/pq"
|
||||
"github.com/mattermost/gorp"
|
||||
"github.com/pkg/errors"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
@@ -186,12 +185,12 @@ func TestStoreLicenseRace(t *testing.T) {
|
||||
}()
|
||||
|
||||
go func() {
|
||||
store.GetReplica()
|
||||
store.GetReplicaX()
|
||||
wg.Done()
|
||||
}()
|
||||
|
||||
go func() {
|
||||
store.GetSearchReplica()
|
||||
store.GetSearchReplicaX()
|
||||
wg.Done()
|
||||
}()
|
||||
|
||||
@@ -276,14 +275,14 @@ func TestGetReplica(t *testing.T) {
|
||||
|
||||
store.UpdateLicense(&model.License{})
|
||||
|
||||
replicas := make(map[*gorp.DbMap]bool)
|
||||
replicas := make(map[*sqlxDBWrapper]bool)
|
||||
for i := 0; i < 5; i++ {
|
||||
replicas[store.GetReplica()] = true
|
||||
replicas[store.GetReplicaX()] = true
|
||||
}
|
||||
|
||||
searchReplicas := make(map[*gorp.DbMap]bool)
|
||||
searchReplicas := make(map[*sqlxDBWrapper]bool)
|
||||
for i := 0; i < 5; i++ {
|
||||
searchReplicas[store.GetSearchReplica()] = true
|
||||
searchReplicas[store.GetSearchReplicaX()] = true
|
||||
}
|
||||
|
||||
if testCase.DataSourceReplicaNum > 0 {
|
||||
@@ -291,13 +290,13 @@ func TestGetReplica(t *testing.T) {
|
||||
assert.Len(t, replicas, testCase.DataSourceReplicaNum)
|
||||
|
||||
for replica := range replicas {
|
||||
assert.NotSame(t, store.GetMaster(), replica)
|
||||
assert.NotSame(t, store.GetMasterX(), replica)
|
||||
}
|
||||
|
||||
} else if assert.Len(t, replicas, 1) {
|
||||
// Otherwise ensure the replicas contains only the master.
|
||||
for replica := range replicas {
|
||||
assert.Same(t, store.GetMaster(), replica)
|
||||
assert.Same(t, store.GetMasterX(), replica)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -306,7 +305,7 @@ func TestGetReplica(t *testing.T) {
|
||||
assert.Len(t, searchReplicas, testCase.DataSourceSearchReplicaNum)
|
||||
|
||||
for searchReplica := range searchReplicas {
|
||||
assert.NotSame(t, store.GetMaster(), searchReplica)
|
||||
assert.NotSame(t, store.GetMasterX(), searchReplica)
|
||||
for replica := range replicas {
|
||||
assert.NotSame(t, searchReplica, replica)
|
||||
}
|
||||
@@ -319,7 +318,7 @@ func TestGetReplica(t *testing.T) {
|
||||
} else if testCase.DataSourceReplicaNum == 0 && assert.Len(t, searchReplicas, 1) {
|
||||
// Otherwise ensure the search replicas contains the master.
|
||||
for searchReplica := range searchReplicas {
|
||||
assert.Same(t, store.GetMaster(), searchReplica)
|
||||
assert.Same(t, store.GetMasterX(), searchReplica)
|
||||
}
|
||||
}
|
||||
})
|
||||
@@ -344,14 +343,14 @@ func TestGetReplica(t *testing.T) {
|
||||
storetest.CleanupSqlSettings(settings)
|
||||
}()
|
||||
|
||||
replicas := make(map[*gorp.DbMap]bool)
|
||||
replicas := make(map[*sqlxDBWrapper]bool)
|
||||
for i := 0; i < 5; i++ {
|
||||
replicas[store.GetReplica()] = true
|
||||
replicas[store.GetReplicaX()] = true
|
||||
}
|
||||
|
||||
searchReplicas := make(map[*gorp.DbMap]bool)
|
||||
searchReplicas := make(map[*sqlxDBWrapper]bool)
|
||||
for i := 0; i < 5; i++ {
|
||||
searchReplicas[store.GetSearchReplica()] = true
|
||||
searchReplicas[store.GetSearchReplicaX()] = true
|
||||
}
|
||||
|
||||
if testCase.DataSourceReplicaNum > 0 {
|
||||
@@ -359,13 +358,13 @@ func TestGetReplica(t *testing.T) {
|
||||
assert.Len(t, replicas, 1)
|
||||
|
||||
for replica := range replicas {
|
||||
assert.Same(t, store.GetMaster(), replica)
|
||||
assert.Same(t, store.GetMasterX(), replica)
|
||||
}
|
||||
|
||||
} else if assert.Len(t, replicas, 1) {
|
||||
// Otherwise ensure the replicas contains only the master.
|
||||
for replica := range replicas {
|
||||
assert.Same(t, store.GetMaster(), replica)
|
||||
assert.Same(t, store.GetMasterX(), replica)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -374,7 +373,7 @@ func TestGetReplica(t *testing.T) {
|
||||
assert.Len(t, searchReplicas, 1)
|
||||
|
||||
for searchReplica := range searchReplicas {
|
||||
assert.Same(t, store.GetMaster(), searchReplica)
|
||||
assert.Same(t, store.GetMasterX(), searchReplica)
|
||||
}
|
||||
|
||||
} else if testCase.DataSourceReplicaNum > 0 {
|
||||
@@ -385,7 +384,7 @@ func TestGetReplica(t *testing.T) {
|
||||
} else if assert.Len(t, searchReplicas, 1) {
|
||||
// Otherwise ensure the search replicas contains the master.
|
||||
for searchReplica := range searchReplicas {
|
||||
assert.Same(t, store.GetMaster(), searchReplica)
|
||||
assert.Same(t, store.GetMasterX(), searchReplica)
|
||||
}
|
||||
}
|
||||
})
|
||||
@@ -755,10 +754,10 @@ func TestExecNoTimeout(t *testing.T) {
|
||||
StoreTest(t, func(t *testing.T, ss store.Store) {
|
||||
sqlStore := ss.(*SqlStore)
|
||||
var query string
|
||||
timeout := sqlStore.master.QueryTimeout
|
||||
sqlStore.master.QueryTimeout = 1
|
||||
timeout := sqlStore.masterX.queryTimeout
|
||||
sqlStore.masterX.queryTimeout = 1
|
||||
defer func() {
|
||||
sqlStore.master.QueryTimeout = timeout
|
||||
sqlStore.masterX.queryTimeout = timeout
|
||||
}()
|
||||
if sqlStore.DriverName() == model.DatabaseDriverMysql {
|
||||
query = `SELECT SLEEP(2);`
|
||||
|
||||
@@ -922,7 +922,7 @@ func (s SqlTeamStore) GetMember(ctx context.Context, teamId string, userId strin
|
||||
}
|
||||
|
||||
var dbMember teamMemberWithSchemeRoles
|
||||
err = s.DBFromContext(ctx).SelectOne(&dbMember, queryString, args...)
|
||||
err = s.DBXFromContext(ctx).Get(&dbMember, queryString, args...)
|
||||
if err != nil {
|
||||
if err == sql.ErrNoRows {
|
||||
return nil, store.NewErrNotFound("TeamMember", fmt.Sprintf("teamId=%s, userId=%s", teamId, userId))
|
||||
@@ -1068,7 +1068,7 @@ func (s SqlTeamStore) GetTeamsForUser(ctx context.Context, userId string) ([]*mo
|
||||
}
|
||||
|
||||
dbMembers := teamMemberWithSchemeRolesList{}
|
||||
_, err = s.SqlStore.DBFromContext(ctx).Select(&dbMembers, queryString, args...)
|
||||
err = s.SqlStore.DBXFromContext(ctx).Select(&dbMembers, queryString, args...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "failed to find TeamMembers with userId=%s", userId)
|
||||
}
|
||||
|
||||
@@ -282,27 +282,22 @@ func themeMigrationFailed(err error) {
|
||||
func upgradeDatabaseToVersion33(sqlStore *SqlStore) {
|
||||
if shouldPerformUpgrade(sqlStore, Version320, Version330) {
|
||||
if sqlStore.DoesColumnExist("Users", "ThemeProps") {
|
||||
params := map[string]interface{}{
|
||||
"Category": model.PreferenceCategoryTheme,
|
||||
"Name": "",
|
||||
}
|
||||
|
||||
transaction, err := sqlStore.GetMaster().Begin()
|
||||
transaction, err := sqlStore.GetMasterX().Beginx()
|
||||
if err != nil {
|
||||
themeMigrationFailed(err)
|
||||
}
|
||||
defer finalizeTransaction(transaction)
|
||||
defer finalizeTransactionX(transaction)
|
||||
|
||||
// copy data across
|
||||
if _, err := transaction.ExecNoTimeout(
|
||||
`INSERT INTO
|
||||
Preferences(UserId, Category, Name, Value)
|
||||
SELECT
|
||||
Id, '`+model.PreferenceCategoryTheme+`', '', ThemeProps
|
||||
Id, '` + model.PreferenceCategoryTheme + `', '', ThemeProps
|
||||
FROM
|
||||
Users
|
||||
WHERE
|
||||
Users.ThemeProps != 'null'`, params); err != nil {
|
||||
Users.ThemeProps != 'null'`); err != nil {
|
||||
themeMigrationFailed(err)
|
||||
return
|
||||
}
|
||||
@@ -554,8 +549,8 @@ func upgradeDatabaseToVersion510(sqlStore *SqlStore) {
|
||||
func upgradeDatabaseToVersion511(sqlStore *SqlStore) {
|
||||
if shouldPerformUpgrade(sqlStore, Version5100, Version5110) {
|
||||
// Enforce all teams have an InviteID set
|
||||
var teams []*model.Team
|
||||
if _, err := sqlStore.GetReplica().Select(&teams, "SELECT * FROM Teams WHERE InviteId = ''"); err != nil {
|
||||
teams := []*model.Team{}
|
||||
if err := sqlStore.GetReplicaX().Select(&teams, "SELECT * FROM Teams WHERE InviteId = ''"); err != nil {
|
||||
mlog.Error("Error fetching Teams without InviteID", mlog.Err(err))
|
||||
} else {
|
||||
for _, team := range teams {
|
||||
@@ -701,7 +696,7 @@ func precheckMigrationToVersion528(sqlStore *SqlStore) error {
|
||||
}
|
||||
|
||||
var schemeIDWrong, typeWrong int
|
||||
row := sqlStore.GetMaster().Db.QueryRow(teamsQuery)
|
||||
row := sqlStore.GetMasterX().QueryRow(teamsQuery)
|
||||
if err = row.Scan(&schemeIDWrong, &typeWrong); err != nil && err != sql.ErrNoRows {
|
||||
return err
|
||||
} else if err == nil && schemeIDWrong > 0 {
|
||||
|
||||
@@ -390,7 +390,7 @@ func (us SqlUserStore) GetMany(ctx context.Context, ids []string) ([]*model.User
|
||||
}
|
||||
|
||||
users := []*model.User{}
|
||||
if _, err := us.SqlStore.DBFromContext(ctx).Select(&users, queryString, args...); err != nil {
|
||||
if err := us.SqlStore.DBXFromContext(ctx).Select(&users, queryString, args...); err != nil {
|
||||
return nil, errors.Wrap(err, "users_get_many_select")
|
||||
}
|
||||
|
||||
@@ -403,7 +403,7 @@ func (us SqlUserStore) Get(ctx context.Context, id string) (*model.User, error)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "users_get_tosql")
|
||||
}
|
||||
row := us.SqlStore.DBFromContext(ctx).Db.QueryRow(queryString, args...)
|
||||
row := us.SqlStore.DBXFromContext(ctx).QueryRow(queryString, args...)
|
||||
|
||||
var user model.User
|
||||
var props, notifyProps, timezone []byte
|
||||
@@ -774,7 +774,7 @@ func (us SqlUserStore) GetAllProfilesInChannel(ctx context.Context, channelID st
|
||||
}
|
||||
|
||||
users := []*model.User{}
|
||||
rows, err := us.SqlStore.DBFromContext(ctx).Db.Query(queryString, args...)
|
||||
rows, err := us.SqlStore.DBXFromContext(ctx).Query(queryString, args...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "failed to find Users")
|
||||
}
|
||||
@@ -998,7 +998,7 @@ func (us SqlUserStore) GetProfileByIds(ctx context.Context, userIds []string, op
|
||||
return nil, errors.Wrap(err, "get_profile_by_ids_tosql")
|
||||
}
|
||||
|
||||
if _, err := us.SqlStore.DBFromContext(ctx).Select(&users, queryString, args...); err != nil {
|
||||
if err := us.SqlStore.DBXFromContext(ctx).Select(&users, queryString, args...); err != nil {
|
||||
return nil, errors.Wrap(err, "failed to find Users")
|
||||
}
|
||||
|
||||
@@ -1679,7 +1679,7 @@ func (us SqlUserStore) GetUsersBatchForIndexing(startTime, endTime int64, limit
|
||||
OrderBy("u.CreateAt").
|
||||
Limit(uint64(limit)).
|
||||
ToSql()
|
||||
_, err := us.GetSearchReplica().Select(&users, usersQuery, args...)
|
||||
err := us.GetSearchReplicaX().Select(&users, usersQuery, args...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "failed to find Users")
|
||||
}
|
||||
@@ -1689,7 +1689,7 @@ func (us SqlUserStore) GetUsersBatchForIndexing(startTime, endTime int64, limit
|
||||
userIds = append(userIds, user.Id)
|
||||
}
|
||||
|
||||
var channelMembers []*model.ChannelMember
|
||||
channelMembers := []*model.ChannelMember{}
|
||||
channelMembersQuery, args, _ := us.getQueryBuilder().
|
||||
Select(`
|
||||
cm.ChannelId,
|
||||
@@ -1709,18 +1709,18 @@ func (us SqlUserStore) GetUsersBatchForIndexing(startTime, endTime int64, limit
|
||||
Join("Channels c ON cm.ChannelId = c.Id").
|
||||
Where(sq.Eq{"c.Type": model.ChannelTypeOpen, "cm.UserId": userIds}).
|
||||
ToSql()
|
||||
_, err = us.GetSearchReplica().Select(&channelMembers, channelMembersQuery, args...)
|
||||
err = us.GetSearchReplicaX().Select(&channelMembers, channelMembersQuery, args...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "failed to find ChannelMembers")
|
||||
}
|
||||
|
||||
var teamMembers []*model.TeamMember
|
||||
teamMembers := []*model.TeamMember{}
|
||||
teamMembersQuery, args, _ := us.getQueryBuilder().
|
||||
Select("TeamId, UserId, Roles, DeleteAt, (SchemeGuest IS NOT NULL AND SchemeGuest) as SchemeGuest, SchemeUser, SchemeAdmin").
|
||||
From("TeamMembers").
|
||||
Where(sq.Eq{"UserId": userIds, "DeleteAt": 0}).
|
||||
ToSql()
|
||||
_, err = us.GetSearchReplica().Select(&teamMembers, teamMembersQuery, args...)
|
||||
err = us.GetSearchReplicaX().Select(&teamMembers, teamMembersQuery, args...)
|
||||
if err != nil {
|
||||
return nil, errors.Wrap(err, "failed to find TeamMembers")
|
||||
}
|
||||
|
||||
@@ -17,15 +17,7 @@ type SqlUserTermsOfServiceStore struct {
|
||||
}
|
||||
|
||||
func newSqlUserTermsOfServiceStore(sqlStore *SqlStore) store.UserTermsOfServiceStore {
|
||||
s := SqlUserTermsOfServiceStore{sqlStore}
|
||||
|
||||
for _, db := range sqlStore.GetAllConns() {
|
||||
table := db.AddTableWithName(model.UserTermsOfService{}, "UserTermsOfService").SetKeys(false, "UserId")
|
||||
table.ColMap("UserId").SetMaxSize(26)
|
||||
table.ColMap("TermsOfServiceId").SetMaxSize(26)
|
||||
}
|
||||
|
||||
return s
|
||||
return SqlUserTermsOfServiceStore{sqlStore}
|
||||
}
|
||||
|
||||
func (s SqlUserTermsOfServiceStore) GetByUser(userId string) (*model.UserTermsOfService, error) {
|
||||
|
||||
@@ -9,8 +9,6 @@ import (
|
||||
"strings"
|
||||
"unicode"
|
||||
|
||||
"github.com/mattermost/gorp"
|
||||
|
||||
"github.com/mattermost/mattermost-server/v6/shared/mlog"
|
||||
)
|
||||
|
||||
@@ -47,14 +45,6 @@ func MapStringsToQueryParams(list []string, paramPrefix string) (string, map[str
|
||||
return "(" + keys.String() + ")", params
|
||||
}
|
||||
|
||||
// finalizeTransaction ensures a transaction is closed after use, rolling back if not already committed.
|
||||
func finalizeTransaction(transaction *gorp.Transaction) {
|
||||
// Rollback returns sql.ErrTxDone if the transaction was already closed.
|
||||
if err := transaction.Rollback(); err != nil && err != sql.ErrTxDone {
|
||||
mlog.Error("Failed to rollback transaction", mlog.Err(err))
|
||||
}
|
||||
}
|
||||
|
||||
// finalizeTransactionX ensures a transaction is closed after use, rolling back if not already committed.
|
||||
func finalizeTransactionX(transaction *sqlxTxWrapper) {
|
||||
// Rollback returns sql.ErrTxDone if the transaction was already closed.
|
||||
|
||||
Ссылка в новой задаче
Block a user