[MM-47976] cmd/mattermost/db: add downgrade commands and enable plan saving (#21779)

Этот коммит содержится в:
Ibrahim Serdar Acikgoz
2023-06-12 12:48:50 +03:00
коммит произвёл GitHub
родитель efaa6264cc
Коммит 4546a2eebb
5 изменённых файлов: 424 добавлений и 105 удалений

229
server/channels/store/sqlstore/migrate.go Обычный файл
Просмотреть файл

@@ -0,0 +1,229 @@
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
// See LICENSE.txt for license information.
package sqlstore
import (
"context"
"fmt"
"log"
"path"
"sort"
"strconv"
"sync"
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/v8/channels/db"
"github.com/mattermost/morph"
"github.com/mattermost/morph/drivers"
ms "github.com/mattermost/morph/drivers/mysql"
ps "github.com/mattermost/morph/drivers/postgres"
"github.com/mattermost/morph/models"
mbindata "github.com/mattermost/morph/sources/embedded"
)
type Migrator struct {
engine *morph.Morph
store *SqlStore
}
func NewMigrator(settings model.SqlSettings, dryRun bool) (*Migrator, error) {
ss := &SqlStore{
rrCounter: 0,
srCounter: 0,
settings: &settings,
quitMonitor: make(chan struct{}),
wgMonitor: &sync.WaitGroup{},
}
ss.initConnection()
ver, err := ss.GetDbVersion(true)
if err != nil {
return nil, fmt.Errorf("error while getting DB version: %w", err)
}
ok, err := ss.ensureMinimumDBVersion(ver)
if !ok {
return nil, fmt.Errorf("error while checking DB version: %w", err)
}
err = ss.ensureDatabaseCollation()
if err != nil {
return nil, fmt.Errorf("error while checking DB collation: %w", err)
}
engine, err := ss.initMorph(dryRun)
if err != nil {
return nil, fmt.Errorf("failed to initialize morph: %w", err)
}
return &Migrator{
engine: engine,
store: ss,
}, nil
}
func (m *Migrator) Close() error {
if err := m.engine.Close(); err != nil {
return fmt.Errorf("failed to close morph engine: %w", err)
}
m.store.Close()
return nil
}
func (m *Migrator) GetFileName(plan *models.Plan) (string, error) {
if len(plan.Migrations) == 0 {
return "", fmt.Errorf("plan is empty")
}
to := plan.Migrations[len(plan.Migrations)-1].Version
from, err := m.store.GetDBSchemaVersion()
if err != nil {
return "", err
}
return fmt.Sprintf("migration_plan_%d_%d", from, to), nil
}
func (ss *SqlStore) initMorph(dryRun bool) (*morph.Morph, error) {
assets := db.Assets()
assetsList, err := assets.ReadDir(path.Join("migrations", ss.DriverName()))
if err != nil {
return nil, err
}
assetNamesForDriver := make([]string, len(assetsList))
for i, entry := range assetsList {
assetNamesForDriver[i] = entry.Name()
}
src, err := mbindata.WithInstance(&mbindata.AssetSource{
Names: assetNamesForDriver,
AssetFunc: func(name string) ([]byte, error) {
return assets.ReadFile(path.Join("migrations", ss.DriverName(), name))
},
})
if err != nil {
return nil, err
}
var driver drivers.Driver
switch ss.DriverName() {
case model.DatabaseDriverMysql:
dataSource, rErr := ResetReadTimeout(*ss.settings.DataSource)
if rErr != nil {
mlog.Fatal("Failed to reset read timeout from datasource.", mlog.Err(rErr), mlog.String("src", *ss.settings.DataSource))
return nil, rErr
}
dataSource, err = AppendMultipleStatementsFlag(dataSource)
if err != nil {
return nil, err
}
db, err2 := SetupConnection("master", dataSource, ss.settings, DBPingAttempts)
if err2 != nil {
return nil, err2
}
driver, err = ms.WithInstance(db)
if err != nil {
return nil, err
}
defer db.Close()
case model.DatabaseDriverPostgres:
driver, err = ps.WithInstance(ss.GetMasterX().DB.DB)
default:
err = fmt.Errorf("unsupported database type %s for migration", ss.DriverName())
}
if err != nil {
return nil, err
}
opts := []morph.EngineOption{
morph.WithLogger(log.New(&morphWriter{}, "", log.Lshortfile)),
morph.WithLock("mm-lock-key"),
morph.SetStatementTimeoutInSeconds(*ss.settings.MigrationsStatementTimeoutSeconds),
morph.SetDryRun(dryRun),
}
engine, err := morph.New(context.Background(), driver, src, opts...)
if err != nil {
return nil, err
}
return engine, nil
}
func (ss *SqlStore) migrate(direction migrationDirection, dryRun bool) error {
engine, err := ss.initMorph(dryRun)
if err != nil {
return err
}
defer engine.Close()
switch direction {
case migrationsDirectionDown:
_, err = engine.ApplyDown(-1)
return err
default:
return engine.ApplyAll()
}
}
func (m *Migrator) GeneratePlan(recover bool) (*models.Plan, error) {
diff, err := m.engine.Diff(models.Up)
if err != nil {
return nil, err
}
plan, err := m.engine.GeneratePlan(diff, recover)
if err != nil {
return nil, err
}
return plan, nil
}
// MigrateWithPlan migrates the database to the latest version using the provided plan.
func (m *Migrator) MigrateWithPlan(plan *models.Plan, dryRun bool) error {
return m.engine.ApplyPlan(plan)
}
func (m *Migrator) DowngradeMigrations(dryRun bool, versions ...string) error {
migrations, err := m.engine.Diff(models.Down)
if err != nil {
return err
}
migrationsToDowngrade := make([]*models.Migration, 0, len(versions))
for _, version := range versions {
for _, migration := range migrations {
versionNumber, sErr := strconv.Atoi(version)
if sErr != nil {
return sErr
}
if migration.Version == uint32(versionNumber) {
migrationsToDowngrade = append(migrationsToDowngrade, migration)
}
}
}
sort.Slice(migrationsToDowngrade, func(i, j int) bool {
return migrationsToDowngrade[i].Version > migrationsToDowngrade[j].Version
})
if len(migrationsToDowngrade) != len(versions) {
mlog.Warn("could not match give migration versions, going to downgrade only the migrations those are available.")
}
plan, err := m.engine.GeneratePlan(migrationsToDowngrade, false)
if err != nil {
return err
}
return m.engine.ApplyPlan(plan)
}

Просмотреть файл

@@ -0,0 +1,33 @@
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
// See LICENSE.txt for license information.
package sqlstore
import (
"testing"
"github.com/mattermost/mattermost/server/public/model"
"github.com/stretchr/testify/assert"
)
func TestUpAndDownMigrations(t *testing.T) {
testDrivers := []string{
model.DatabaseDriverPostgres,
model.DatabaseDriverMysql,
}
for _, driver := range testDrivers {
t.Run("Should be reversible for "+driver, func(t *testing.T) {
settings, err := makeSqlSettings(driver)
if err != nil {
t.Skip(err)
}
store := New(*settings, nil)
defer store.Close()
err = store.migrate(migrationsDirectionDown, false)
assert.NoError(t, err, "downing migrations should not error")
})
}
}

Просмотреть файл

@@ -8,31 +8,22 @@ import (
"database/sql"
dbsql "database/sql"
"fmt"
"log"
"path"
"strconv"
"strings"
"sync"
"sync/atomic"
"time"
"github.com/mattermost/morph"
sq "github.com/mattermost/squirrel"
"github.com/mattermost/morph/drivers"
ms "github.com/mattermost/morph/drivers/mysql"
ps "github.com/mattermost/morph/drivers/postgres"
"github.com/go-sql-driver/mysql"
_ "github.com/golang-migrate/migrate/v4/source/file"
"github.com/jmoiron/sqlx"
"github.com/lib/pq"
mbindata "github.com/mattermost/morph/sources/embedded"
"github.com/pkg/errors"
"github.com/mattermost/mattermost/server/public/model"
"github.com/mattermost/mattermost/server/public/shared/mlog"
"github.com/mattermost/mattermost/server/v8/channels/db"
"github.com/mattermost/mattermost/server/v8/channels/store"
"github.com/mattermost/mattermost/server/v8/einterfaces"
)
@@ -177,7 +168,7 @@ func New(settings model.SqlSettings, metrics einterfaces.MetricsInterface) *SqlS
mlog.Fatal("Error while checking DB collation.", mlog.Err(err))
}
err = store.migrate(migrationsDirectionUp)
err = store.migrate(migrationsDirectionUp, false)
if err != nil {
mlog.Fatal("Failed to apply database migrations.", mlog.Err(err))
}
@@ -1191,76 +1182,6 @@ func (ss *SqlStore) hasLicense() bool {
return hasLicense
}
func (ss *SqlStore) migrate(direction migrationDirection) error {
assets := db.Assets()
assetsList, err := assets.ReadDir(path.Join("migrations", ss.DriverName()))
if err != nil {
return err
}
assetNamesForDriver := make([]string, len(assetsList))
for i, entry := range assetsList {
assetNamesForDriver[i] = entry.Name()
}
src, err := mbindata.WithInstance(&mbindata.AssetSource{
Names: assetNamesForDriver,
AssetFunc: func(name string) ([]byte, error) {
return assets.ReadFile(path.Join("migrations", ss.DriverName(), name))
},
})
if err != nil {
return err
}
var driver drivers.Driver
switch ss.DriverName() {
case model.DatabaseDriverMysql:
dataSource, rErr := ResetReadTimeout(*ss.settings.DataSource)
if rErr != nil {
mlog.Fatal("Failed to reset read timeout from datasource.", mlog.Err(rErr), mlog.String("src", *ss.settings.DataSource))
return rErr
}
dataSource, err = AppendMultipleStatementsFlag(dataSource)
if err != nil {
return err
}
db, err2 := SetupConnection("master", dataSource, ss.settings, DBPingAttempts)
if err2 != nil {
return err2
}
driver, err = ms.WithInstance(db)
defer db.Close()
case model.DatabaseDriverPostgres:
driver, err = ps.WithInstance(ss.GetMasterX().DB.DB)
default:
err = fmt.Errorf("unsupported database type %s for migration", ss.DriverName())
}
if err != nil {
return err
}
opts := []morph.EngineOption{
morph.WithLogger(log.New(&morphWriter{}, "", log.Lshortfile)),
morph.WithLock("mm-lock-key"),
morph.SetStatementTimeoutInSeconds(*ss.settings.MigrationsStatementTimeoutSeconds),
}
engine, err := morph.New(context.Background(), driver, src, opts...)
if err != nil {
return err
}
defer engine.Close()
switch direction {
case migrationsDirectionDown:
_, err = engine.ApplyDown(-1)
return err
default:
return engine.ApplyAll()
}
}
func convertMySQLFullTextColumnsToPostgres(columnNames string) string {
columns := strings.Split(columnNames, ", ")
concatenatedColumnNames := ""

Просмотреть файл

@@ -558,27 +558,6 @@ func TestIsBinaryParamEnabled(t *testing.T) {
}
func TestUpAndDownMigrations(t *testing.T) {
testDrivers := []string{
model.DatabaseDriverPostgres,
model.DatabaseDriverMysql,
}
for _, driver := range testDrivers {
t.Run("Should be reversible for "+driver, func(t *testing.T) {
settings, err := makeSqlSettings(driver)
if err != nil {
t.Skip(err)
}
store := New(*settings, nil)
defer store.Close()
err = store.migrate(migrationsDirectionDown)
assert.NoError(t, err, "downing migrations should not error")
})
}
}
func TestGetAllConns(t *testing.T) {
t.Parallel()
testCases := []struct {
@@ -777,7 +756,7 @@ func TestReplicaLagQuery(t *testing.T) {
require.NoError(t, store.initConnection())
store.stores.post = newSqlPostStore(store, mockMetrics)
err = store.migrate(migrationsDirectionUp)
err = store.migrate(migrationsDirectionUp, false)
require.NoError(t, err)
defer store.Close()