Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
79 changes: 65 additions & 14 deletions go/logic/applier.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ import (

"github.com/github/gh-ost/go/mysql"
drivermysql "github.com/go-sql-driver/mysql"
version "github.com/hashicorp/go-version"
"github.com/openark/golib/sqlutils"
)

Expand Down Expand Up @@ -92,6 +93,10 @@ type Applier struct {
migrationLockName string
migrationLockStop chan struct{}
migrationLockDone chan struct{}

// usePerformanceSchemaProcesslist tells ExpectProcess to read performance_schema.processlist
// rather than information_schema.processlist; see chooseProcesslistTable()
usePerformanceSchemaProcesslist bool
}

func NewApplier(migrationContext *base.MigrationContext) *Applier {
Expand Down Expand Up @@ -140,6 +145,7 @@ func (apl *Applier) InitDBConnections() (err error) {
if err := apl.validateAndReadGlobalVariables(); err != nil {
return err
}
apl.chooseProcesslistTable()
if !apl.migrationContext.AliyunRDS && !apl.migrationContext.GoogleCloudPlatform && !apl.migrationContext.AzureMySQL {
if impliedKey, err := mysql.GetInstanceKey(apl.db); err != nil {
return err
Expand Down Expand Up @@ -351,6 +357,47 @@ func (apl *Applier) validateAndReadGlobalVariables() error {
return nil
}

// chooseProcesslistTable decides which table ExpectProcess reads. performance_schema.processlist is
// preferred: since MySQL 9.1, reading information_schema.processlist fails with Error 3854 while any
// other session runs a statement not representable in utf8mb3, such as one with an emoji.
// information_schema.processlist remains the fallback when performance_schema.processlist is not usable:
// absent (MySQL 5.7, 8.0 before 8.0.22, MariaDB), empty because performance_schema is disabled, or not
// readable by this account.
func (apl *Applier) chooseProcesslistTable() {
var ownSessions int64
query := `select /* gh-ost */ count(*) from performance_schema.processlist where id = connection_id()`
err := apl.db.QueryRow(query).Scan(&ownSessions)
if err == nil && ownSessions == 1 {
apl.usePerformanceSchemaProcesslist = true
apl.migrationContext.Log.Infof("will use performance_schema.processlist on applier")
return
}
reason := "own session not listed, performance_schema may be disabled"
if err != nil {
reason = err.Error()
}
if hasStrictTmpTableCharsetConversion(apl.migrationContext.ApplierMySQLVersion) {
apl.migrationContext.Log.Warningf("performance_schema.processlist is not usable on applier (%s); falling back to information_schema.processlist. On MySQL 9.1 and later reading it fails while any other session runs a statement not representable in utf8mb3 (such as one with an emoji), and the atomic cut-over may then fail to find its RENAME. Enable performance_schema, and make sure this user can read performance_schema.processlist, to avoid this. See https://github.com/github/gh-ost/issues/1780 for details", reason)
return
}
apl.migrationContext.Log.Infof("will use information_schema.processlist on applier: performance_schema.processlist is not usable (%s)", reason)
}

// hasStrictTmpTableCharsetConversion tells whether the server is MySQL 9.1 or later, where character
// set conversions into temporary tables are strict: information_schema.processlist then fails with
// Error 3854 instead of a warning while a session runs a statement not representable in utf8mb3
// (MySQL Bug #87579).
func hasStrictTmpTableCharsetConversion(mysqlVersion string) bool {
if mysql.IsMariaDB(mysqlVersion) {
return false
}
serverVersion, err := version.NewVersion(mysqlVersion)
if err != nil {
return false
}
return serverVersion.Core().GreaterThanOrEqual(version.Must(version.NewVersion("9.1")))
}

// generateSqlModeQuery return a `sql_mode = ...` query, to be wrapped with a `set session` or `set global`,
// based on gh-ost configuration:
// - User may skip strict mode
Expand Down Expand Up @@ -1443,28 +1490,32 @@ func (apl *Applier) ExpectUsedLock(sessionId int64) error {
return nil
}

// ExpectProcess expects a process to show up in `SHOW PROCESSLIST` that has given characteristics
// ExpectProcess expects a process to show up in `SHOW PROCESSLIST` that has given characteristics.
// It reads performance_schema.processlist when usable, information_schema.processlist otherwise.
func (apl *Applier) ExpectProcess(sessionId int64, stateHint, infoHint string) error {
found := false
query := `
processlistTable := "information_schema.processlist"
if apl.usePerformanceSchemaProcesslist {
processlistTable = "performance_schema.processlist"
}
query := fmt.Sprintf(`
select /* gh-ost */ id
from
information_schema.processlist
%s
where
id != connection_id()
and ? in (0, id)
and state like concat('%', ?, '%')
and info like concat('%', ?, '%')`
err := sqlutils.QueryRowsMap(apl.db, query, func(m sqlutils.RowMap) error {
found = true
return nil
}, sessionId, stateHint, infoHint)
if err != nil {
return err
}
if !found {
and state like concat('%%', ?, '%%')
and info like concat('%%', ?, '%%')`, processlistTable)
// A failed read may only report its error after the result set has started (Error 3854 does);
// sqlutils.QueryRowsMap would drop it and report the process as missing, Scan() returns it.
var id int64
err := apl.db.QueryRow(query, sessionId, stateHint, infoHint).Scan(&id)
if errors.Is(err, gosql.ErrNoRows) {
return fmt.Errorf("cannot find process. Hints: %s, %s", stateHint, infoHint)
}
if err != nil {
return apl.migrationContext.Log.Errore(fmt.Errorf("failed reading %s: %w", processlistTable, err))
}
return nil
}

Expand Down
222 changes: 222 additions & 0 deletions go/logic/applier_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
gosql "database/sql"
"errors"
"strings"
"sync"
"testing"
"time"

Expand Down Expand Up @@ -333,6 +334,78 @@ func TestClassifyAnalyzeTableResult(t *testing.T) {
}
}

func TestHasStrictTmpTableCharsetConversion(t *testing.T) {
tests := []struct {
version string
strict bool
}{
{version: "5.7.41-log", strict: false},
{version: "8.0.42", strict: false},
{version: "8.0.41-32", strict: false},
{version: "8.4.3", strict: false},
{version: "9.0.1", strict: false},
{version: "9.1.0", strict: true},
{version: "9.1.0-commercial", strict: true},
{version: "9.7.1-1", strict: true},
{version: "9.7.2", strict: true},
{version: "10.6.27-MariaDB-log", strict: false},
{version: "11.8.8-MariaDB-ubu2404", strict: false},
{version: "", strict: false},
}
for _, tc := range tests {
t.Run(tc.version, func(t *testing.T) {
require.Equal(t, tc.strict, hasStrictTmpTableCharsetConversion(tc.version))
})
}
}

// lockTestTableAndQueue write-locks the test table, then runs each of the given statements on a
// connection of its own and waits until it is blocked on a metadata lock behind that table lock.
// It returns the session ids of these connections, in order, and a function that releases the table
// lock and waits for the statements to finish; their outcome is not checked.
func lockTestTableAndQueue(t *testing.T, db *gosql.DB, statements ...string) (sessionIds []int64, unlock func()) {
ctx := context.Background()

lockConn, err := db.Conn(ctx)
require.NoError(t, err)
var wg sync.WaitGroup
unlock = func() {
_, _ = lockConn.ExecContext(ctx, "UNLOCK TABLES")
_ = lockConn.Close()
wg.Wait()
}
queued := false
defer func() {
if !queued {
unlock()
}
}()

_, err = lockConn.ExecContext(ctx, "LOCK TABLES "+getTestTableName()+" WRITE")
require.NoError(t, err)

for _, statement := range statements {
conn, err := db.Conn(ctx)
require.NoError(t, err)
var sessionId int64
require.NoError(t, conn.QueryRowContext(ctx, "SELECT CONNECTION_ID()").Scan(&sessionId))
wg.Add(1)
go func() {
defer wg.Done()
defer conn.Close()
_, _ = conn.ExecContext(ctx, statement)
}()
require.Eventually(t, func() bool {
var waiting int64
err := db.QueryRowContext(ctx, `SELECT COUNT(*) FROM performance_schema.processlist WHERE id = ? AND state LIKE '%metadata lock%'`, sessionId).Scan(&waiting)
return err == nil && waiting == 1
}, 10*time.Second, 10*time.Millisecond, "statement is not waiting for the table lock: %s", statement)
sessionIds = append(sessionIds, sessionId)
}
queued = true
return sessionIds, unlock
}

type ApplierTestSuite struct {
suite.Suite

Expand Down Expand Up @@ -407,6 +480,7 @@ func (suite *ApplierTestSuite) TestInitDBConnections() {
suite.Require().Equal(mysqlVersion, migrationContext.ApplierMySQLVersion)
suite.Require().Equal(int64(28800), migrationContext.ApplierWaitTimeout)
suite.Require().Equal("+00:00", migrationContext.ApplierTimeZone)
suite.Require().True(applier.usePerformanceSchemaProcesslist)

suite.Require().Equal(sql.NewColumnList([]string{"id", "item_id"}), migrationContext.OriginalTableColumnsOnApplier)
}
Expand Down Expand Up @@ -631,6 +705,44 @@ func (suite *ApplierTestSuite) TestAcquireMigrationLockFailsWhenHeld() {
suite.Require().Nil(applierB.migrationLockConn)
}

func (suite *ApplierTestSuite) TestExpectProcess() {
ctx := context.Background()

_, err := suite.db.ExecContext(ctx, fmt.Sprintf("CREATE TABLE %s (id INT PRIMARY KEY, name VARCHAR(64))", getTestTableName()))
suite.Require().NoError(err)

connectionConfig, err := getTestConnectionConfig(ctx, suite.mysqlContainer)
suite.Require().NoError(err)

migrationContext := newTestMigrationContext()
migrationContext.ApplierConnectionConfig = connectionConfig
migrationContext.SetConnectionConfig("innodb")

applier := NewApplier(migrationContext)
defer applier.Teardown()
suite.Require().NoError(applier.InitDBConnections())
suite.Require().True(applier.usePerformanceSchemaProcesslist)

sessionIds, unlock := lockTestTableAndQueue(suite.T(), suite.db,
fmt.Sprintf("RENAME TABLE %s TO %s", getTestTableName(), getTestOldTableName()),
)
defer func() {
unlock()
_, err := suite.db.ExecContext(ctx, "DROP TABLE IF EXISTS "+getTestOldTableName())
suite.Require().NoError(err)
}()
renameSessionId := sessionIds[0]

suite.Require().NoError(applier.ExpectProcess(renameSessionId, "metadata lock", "rename"))
suite.Require().NoError(applier.ExpectProcess(0, "metadata lock", "rename"))
suite.Require().ErrorContains(applier.ExpectProcess(renameSessionId, "metadata lock", "insert"), "cannot find process")

// information_schema.processlist, the fallback, finds the same session
applier.usePerformanceSchemaProcesslist = false
suite.Require().NoError(applier.ExpectProcess(renameSessionId, "metadata lock", "rename"))
suite.Require().ErrorContains(applier.ExpectProcess(renameSessionId, "metadata lock", "insert"), "cannot find process")
}

func TestBuildMigrationLockName(t *testing.T) {
t.Run("short name is returned verbatim", func(t *testing.T) {
name := buildMigrationLockName("mydb", "mytable")
Expand Down Expand Up @@ -1734,3 +1846,113 @@ func TestApplier(t *testing.T) {
}
suite.Run(t, new(ApplierTestSuite))
}

// testMysql9ContainerImage is a MySQL release with strict character set conversions into temporary
// tables (9.1 and later)
const testMysql9ContainerImage = "mysql:9.7.2"

// ApplierMySQL9TestSuite runs against MySQL 9.1 or later, where character set conversions into
// temporary tables are strict: information_schema.processlist cannot be read while any session runs
// a statement not representable in utf8mb3.
type ApplierMySQL9TestSuite struct {
suite.Suite

mysqlContainer testcontainers.Container
db *gosql.DB
}

func (suite *ApplierMySQL9TestSuite) SetupSuite() {
ctx := context.Background()
// my.cnf.test is not used: MySQL 9.7 rejects innodb_log_file_size
mysqlContainer, err := testmysql.Run(ctx,
testMysql9ContainerImage,
testmysql.WithDatabase(testMysqlDatabase),
testmysql.WithUsername(testMysqlUser),
testmysql.WithPassword(testMysqlPass),
)
suite.Require().NoError(err)

suite.mysqlContainer = mysqlContainer

dsn, err := mysqlContainer.ConnectionString(ctx)
suite.Require().NoError(err)

db, err := gosql.Open("mysql", dsn)
suite.Require().NoError(err)

suite.db = db
}

func (suite *ApplierMySQL9TestSuite) TearDownSuite() {
suite.Assert().NoError(suite.db.Close())
suite.Assert().NoError(testcontainers.TerminateContainer(suite.mysqlContainer))
}

func (suite *ApplierMySQL9TestSuite) TearDownTest() {
ctx := context.Background()

_, err := suite.db.ExecContext(ctx, "DROP TABLE IF EXISTS "+getTestTableName())
suite.Require().NoError(err)
_, err = suite.db.ExecContext(ctx, "DROP TABLE IF EXISTS "+getTestOldTableName())
suite.Require().NoError(err)
}

func (suite *ApplierMySQL9TestSuite) TestExpectProcessWhileUtf8mb4StatementIsRunning() {
ctx := context.Background()

_, err := suite.db.ExecContext(ctx, fmt.Sprintf("CREATE TABLE %s (id INT PRIMARY KEY, name VARCHAR(64))", getTestTableName()))
suite.Require().NoError(err)

connectionConfig, err := getTestConnectionConfig(ctx, suite.mysqlContainer)
suite.Require().NoError(err)

migrationContext := newTestMigrationContext()
migrationContext.ApplierConnectionConfig = connectionConfig
migrationContext.SetConnectionConfig("innodb")

applier := NewApplier(migrationContext)
defer applier.Teardown()
suite.Require().NoError(applier.InitDBConnections())
suite.Require().True(applier.usePerformanceSchemaProcesslist)

// A write carrying an emoji waits behind the table lock, as writes do behind the lock of the
// atomic cut-over, and so does the RENAME after it.
sessionIds, unlock := lockTestTableAndQueue(suite.T(), suite.db,
fmt.Sprintf("INSERT INTO %s (id, name) VALUES (1, '\U0001F600')", getTestTableName()),
fmt.Sprintf("RENAME TABLE %s TO %s", getTestTableName(), getTestOldTableName()),
)
defer unlock()
renameSessionId := sessionIds[1]

// Control: information_schema.processlist cannot be read on this server meanwhile.
var id int64
err = suite.db.QueryRowContext(ctx, `
select id
from
information_schema.processlist
where
id != connection_id()
and ? in (0, id)
and state like concat('%', ?, '%')
and info like concat('%', ?, '%')`,
renameSessionId, "metadata lock", "rename",
).Scan(&id)
var mysqlErr *drivermysql.MySQLError
suite.Require().ErrorAs(err, &mysqlErr)
suite.Require().Equal(uint16(3854), mysqlErr.Number) // ER_CANNOT_CONVERT_STRING

suite.Require().NoError(applier.ExpectProcess(renameSessionId, "metadata lock", "rename"))

// The information_schema.processlist fallback reports the failed read, rather than a missing process.
applier.usePerformanceSchemaProcesslist = false
err = applier.ExpectProcess(renameSessionId, "metadata lock", "rename")
suite.Require().ErrorAs(err, &mysqlErr)
suite.Require().Equal(uint16(3854), mysqlErr.Number)
}

func TestApplierMySQL9(t *testing.T) {
if testing.Short() {
t.Skip("skipping MySQL 9 applier test suite in short mode")
}
suite.Run(t, new(ApplierMySQL9TestSuite))
}