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
10 changes: 6 additions & 4 deletions DEPS.bzl
Original file line number Diff line number Diff line change
Expand Up @@ -3777,8 +3777,9 @@ def go_deps():
name = "com_github_pingcap_kvproto",
build_file_proto_mode = "disable_global",
importpath = "github.com/pingcap/kvproto",
sum = "h1:qW0gHsqY3X3qmyiEfESqpYSF3Vu7agUiAyBnPWQXtm8=",
version = "v0.0.0-20260721064811-683dad8fa368",
replace = "github.com/wfxr/kvproto",
sum = "h1:wQuil8SCJhSp+LJqqiMgPjzBZgzobJCIQFp52UIYq6Y=",
version = "v0.0.0-20260806092442-d04fa0402753",
)
go_repository(
name = "com_github_pingcap_log",
Expand Down Expand Up @@ -4492,8 +4493,9 @@ def go_deps():
build_tags = ["nextgen", "intest"],
build_file_proto_mode = "disable_global",
importpath = "github.com/tikv/client-go/v2",
sum = "h1:20uV3D/EvkPVkDvSGWcOg2+jZOKaCYpIoIkwSb7amJM=",
version = "v2.0.8-0.20260803074519-341d4692ec57",
replace = "github.com/wfxr/client-go/v2",
sum = "h1:mqjqOtWPO+TVbH8lxEghT2DBcoZqR5vsQQLVwUKJgYo=",
version = "v2.0.8-0.20260812071238-dcfbdfcfb5af",
)
go_repository(
name = "com_github_tikv_pd_client",
Expand Down
6 changes: 4 additions & 2 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,7 @@ require (
github.com/pingcap/errors v0.11.5-0.20260508054701-306e305bcf41
github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86
github.com/pingcap/fn v1.0.0
github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368
github.com/pingcap/kvproto v0.0.0-20260806092442-d04fa0402753
github.com/pingcap/log v1.1.1-0.20250917021125-19901e015dc9
github.com/pingcap/metering_sdk v0.0.0-20260324055927-14fead745f1d
github.com/pingcap/sysutil v1.0.1-0.20240311050922-ae81ee01f3a5
Expand All @@ -121,7 +121,7 @@ require (
github.com/stathat/consistent v1.0.0
github.com/stretchr/testify v1.11.1
github.com/tiancaiamao/appdash v0.0.0-20181126055449-889f96f722a2
github.com/tikv/client-go/v2 v2.0.8-0.20260803074519-341d4692ec57
github.com/tikv/client-go/v2 v2.0.8-0.20260812071238-dcfbdfcfb5af
github.com/tikv/pd/client v0.0.0-20260720043438-0b37df9a48ed
github.com/timakin/bodyclose v0.0.0-20241222091800-1db5c5ca4d67
github.com/twmb/murmur3 v1.1.6
Expand Down Expand Up @@ -363,7 +363,9 @@ require (
replace (
cloud.google.com/go/storage => cloud.google.com/go/storage v1.39.1
github.com/go-ldap/ldap/v3 => github.com/YangKeao/ldap/v3 v3.4.5-0.20230421065457-369a3bab1117
github.com/pingcap/kvproto => github.com/wfxr/kvproto v0.0.0-20260806092442-d04fa0402753
github.com/pingcap/tidb/pkg/parser => ./pkg/parser
github.com/tikv/client-go/v2 => github.com/wfxr/client-go/v2 v2.0.8-0.20260812071238-dcfbdfcfb5af
// TODO: `sourcegraph.com/sourcegraph/appdash` has been archived, and the original host has been removed.
// Please remove these dependencies.
sourcegraph.com/sourcegraph/appdash => github.com/sourcegraph/appdash v0.0.0-20190731080439-ebfcffb1b5c0
Expand Down
95 changes: 85 additions & 10 deletions go.sum

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions pkg/errno/errcode.go
Original file line number Diff line number Diff line change
Expand Up @@ -1197,5 +1197,6 @@ const (
ErrTiFlashServerTimeout = 9012
ErrTiFlashServerBusy = 9013
ErrTiFlashBackfillIndex = 9014
ErrSharedLockLost = 9015
ErrUserPrefixMismatch = 20003
)
8 changes: 6 additions & 2 deletions pkg/errno/errname.go
Original file line number Diff line number Diff line change
Expand Up @@ -1184,8 +1184,12 @@ var MySQLErrName = map[uint16]*mysql.ErrMessage{
ErrTiFlashServerTimeout: mysql.Message("TiFlash server timeout", nil),
ErrTiFlashServerBusy: mysql.Message("TiFlash server is busy", nil),
ErrTiFlashBackfillIndex: mysql.Message("TiFlash backfill index failed: %s", nil),
ErrResolveLockTimeout: mysql.Message("Resolve lock timeout", nil),
ErrRegionUnavailable: mysql.Message("Region is unavailable", nil),
ErrSharedLockLost: mysql.Message(
"Shared lock was lost during lock upgrade; transaction cannot continue, txnStartTS=%d, key=%s",
[]int{1},
),
ErrResolveLockTimeout: mysql.Message("Resolve lock timeout", nil),
ErrRegionUnavailable: mysql.Message("Region is unavailable", nil),
// In most cases, the error `ErrTxnAbortedByGC` is caused by the transaction runs too long, instead of improper GC
// life time configuration. This means the description of this error is not accurate.
// However, as this error message is already widely acknowledged and might have become part of our diagnosing
Expand Down
1 change: 1 addition & 0 deletions pkg/executor/select.go
Original file line number Diff line number Diff line change
Expand Up @@ -345,6 +345,7 @@ func newLockCtx(sctx sessionctx.Context, lockWaitTime int64, numKeys int, inShar
lockCtx.Killed = &seVars.SQLKiller.Signal
lockCtx.LockExpired = &seVars.TxnCtx.LockExpire
lockCtx.InShareMode = inSharedMode
lockCtx.AllowSharedLockUpgrade = seVars.EnableSharedLockUpgrade

// Set max_execution_time deadline for SELECT statements
maxExectionTime := seVars.GetMaxExecutionTime()
Expand Down
89 changes: 84 additions & 5 deletions pkg/executor/select_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,26 +12,86 @@
// See the License for the specific language governing permissions and
// limitations under the License.

package executor_test
package executor

import (
"context"
"fmt"
"testing"

"github.com/pingcap/tidb/pkg/domain"
"github.com/pingcap/tidb/pkg/executor"
"github.com/pingcap/tidb/pkg/infoschema"
"github.com/pingcap/tidb/pkg/kv"
"github.com/pingcap/tidb/pkg/parser/ast"
"github.com/pingcap/tidb/pkg/parser/mysql"
"github.com/pingcap/tidb/pkg/sessionctx"
"github.com/pingcap/tidb/pkg/sessiontxn"
"github.com/pingcap/tidb/pkg/util/mock"
"github.com/stretchr/testify/require"
)

type stubTxnManager struct {
forUpdateTS uint64
}

func (m stubTxnManager) AdviseWarmup() error { return nil }

func (m stubTxnManager) AdviseOptimizeWithPlan(any) error { return nil }

func (m stubTxnManager) GetTxnInfoSchema() infoschema.InfoSchema { return nil }

func (m stubTxnManager) GetTxnScope() string { return "" }

func (m stubTxnManager) GetReadReplicaScope() string { return "" }

func (m stubTxnManager) GetStmtReadTS() (uint64, error) { return 0, nil }

func (m stubTxnManager) GetStmtForUpdateTS() (uint64, error) { return m.forUpdateTS, nil }

func (m stubTxnManager) GetContextProvider() sessiontxn.TxnContextProvider { return nil }

func (m stubTxnManager) GetSnapshotWithStmtReadTS() (kv.Snapshot, error) { return nil, nil }

func (m stubTxnManager) GetSnapshotWithStmtForUpdateTS() (kv.Snapshot, error) { return nil, nil }

func (m stubTxnManager) EnterNewTxn(context.Context, *sessiontxn.EnterNewTxnRequest) error {
return nil
}

func (m stubTxnManager) OnTxnEnd() {}

func (m stubTxnManager) OnStmtStart(context.Context, ast.StmtNode) error { return nil }

func (m stubTxnManager) OnPessimisticStmtStart(context.Context) error { return nil }

func (m stubTxnManager) OnPessimisticStmtEnd(context.Context, bool) error { return nil }

func (m stubTxnManager) OnStmtErrorForNextAction(context.Context, sessiontxn.StmtErrorHandlePoint, error) (sessiontxn.StmtErrorAction, error) {
return sessiontxn.StmtActionNoIdea, nil
}

func (m stubTxnManager) OnStmtRetry(context.Context) error { return nil }

func (m stubTxnManager) OnStmtCommit(context.Context) error { return nil }

func (m stubTxnManager) OnStmtRollback(context.Context, bool) error { return nil }

func (m stubTxnManager) OnStmtEnd() {}

func (m stubTxnManager) OnLocalTemporaryTableCreated() {}

func (m stubTxnManager) ActivateTxn() (kv.Transaction, error) { return nil, nil }

func (m stubTxnManager) GetCurrentStmt() ast.StmtNode { return nil }

func (m stubTxnManager) SetOptionsBeforeCommit(kv.Transaction, func(uint64) bool) error { return nil }

func BenchmarkResetContextOfStmt(b *testing.B) {
stmt := &ast.SelectStmt{}
ctx := mock.NewContext()
ctx.BindDomainAndSchValidator(&domain.Domain{}, nil)
for i := 0; i < b.N; i++ {
executor.ResetContextOfStmt(ctx, stmt)
ResetContextOfStmt(ctx, stmt)
}
}

Expand All @@ -58,13 +118,32 @@ func TestImportIntoShouldHaveSameFlagsAsInsert(t *testing.T) {
mode, err := mysql.GetSQLMode(modeStr)
require.NoError(t, err)
insertCtx.GetSessionVars().SQLMode = mode
require.NoError(t, executor.ResetContextOfStmt(insertCtx, insertStmt))
require.NoError(t, ResetContextOfStmt(insertCtx, insertStmt))
importCtx.GetSessionVars().SQLMode = mode
require.NoError(t, executor.ResetContextOfStmt(importCtx, importStmt))
require.NoError(t, ResetContextOfStmt(importCtx, importStmt))

insertTypeCtx := insertCtx.GetSessionVars().StmtCtx.TypeCtx()
importTypeCtx := importCtx.GetSessionVars().StmtCtx.TypeCtx()
require.EqualValues(t, insertTypeCtx.Flags(), importTypeCtx.Flags())
})
}

t.Run("shared lock upgrade gate propagates to lock ctx", func(t *testing.T) {
originalGetTxnManager := sessiontxn.GetTxnManager
sessiontxn.GetTxnManager = func(sctx sessionctx.Context) sessiontxn.TxnManager {
return stubTxnManager{forUpdateTS: 9527}
}
t.Cleanup(func() {
sessiontxn.GetTxnManager = originalGetTxnManager
})

sctx := mock.NewContext()
sctx.GetSessionVars().EnableSharedLockUpgrade = true

lockCtx, err := newLockCtx(sctx, 123, 1, true)
require.NoError(t, err)
require.True(t, lockCtx.InShareMode)
require.True(t, lockCtx.AllowSharedLockUpgrade)
require.Equal(t, uint64(9527), lockCtx.ForUpdateTS)
})
}
2 changes: 2 additions & 0 deletions pkg/kv/error.go
Original file line number Diff line number Diff line change
Expand Up @@ -70,6 +70,8 @@ var (
mysql.MySQLErrName[mysql.ErrWriteConflictInTiDB].RedactArgPos,
),
)
// ErrSharedLockLost means a shared lock was confirmed lost during upgrade.
ErrSharedLockLost = dbterror.ClassTiKV.NewStd(mysql.ErrSharedLockLost)
// ErrLockExpire is the error when the lock is expired.
ErrLockExpire = dbterror.ClassTiKV.NewStd(mysql.ErrLockExpire)
// ErrAssertionFailed is the error when an assertion fails.
Expand Down
1 change: 1 addition & 0 deletions pkg/kv/error_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ func TestError(t *testing.T) {
ErrNotImplemented,
ErrWriteConflict,
ErrWriteConflictInTiDB,
ErrSharedLockLost,
}

for _, err := range kvErrs {
Expand Down
3 changes: 3 additions & 0 deletions pkg/session/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -196,11 +196,13 @@ go_test(
"//pkg/parser/auth",
"//pkg/parser/format",
"//pkg/parser/mysql",
"//pkg/session/metrics",
"//pkg/session/sessionapi",
"//pkg/sessionctx",
"//pkg/sessionctx/stmtctx",
"//pkg/sessionctx/vardef",
"//pkg/sessionctx/variable",
"//pkg/sessiontxn",
"//pkg/statistics",
"//pkg/store",
"//pkg/store/mockstore",
Expand All @@ -214,6 +216,7 @@ go_test(
"//pkg/util",
"//pkg/util/benchdaily",
"//pkg/util/chunk",
"//pkg/util/dbterror/exeerrors",
"//pkg/util/execdetails",
"//pkg/util/logutil",
"//pkg/util/memory",
Expand Down
22 changes: 20 additions & 2 deletions pkg/session/tidb.go
Original file line number Diff line number Diff line change
Expand Up @@ -298,6 +298,16 @@ func shouldCheckConnectionAliveBeforeCommit(sessVars *variable.SessionVars, sql
}
}

func shouldRollbackTxnOnError(txn kv.Transaction, err error) bool {
if !txn.Valid() {
return false
}
if kv.ErrSharedLockLost.Equal(err) {
return true
}
return txn.IsPessimistic() && exeerrors.ErrDeadlock.Equal(err)
}

func autoCommitAfterStmt(ctx context.Context, se *session, meetsErr error, sql sqlexec.Statement) error {
isInternal := false
if internal := se.txn.GetOption(kv.RequestSourceInternal); internal != nil && internal.(bool) {
Expand All @@ -309,8 +319,16 @@ func autoCommitAfterStmt(ctx context.Context, se *session, meetsErr error, sql s
logutil.BgLogger().Info("rollbackTxn called due to ddl/autocommit failure")
se.RollbackTxn(ctx)
recordAbortTxnDuration(sessVars, isInternal)
} else if se.txn.Valid() && se.txn.IsPessimistic() && exeerrors.ErrDeadlock.Equal(meetsErr) {
logutil.BgLogger().Info("rollbackTxn for deadlock", zap.Uint64("txn", se.txn.StartTS()))
} else if shouldRollbackTxnOnError(&se.txn, meetsErr) {
if kv.ErrSharedLockLost.Equal(meetsErr) {
logutil.BgLogger().Info(
"rollbackTxn for shared lock loss",
zap.Uint64("txn", se.txn.StartTS()),
zap.Error(meetsErr),
)
} else {
logutil.BgLogger().Info("rollbackTxn for deadlock", zap.Uint64("txn", se.txn.StartTS()))
}
se.RollbackTxn(ctx)
recordAbortTxnDuration(sessVars, isInternal)
}
Expand Down
94 changes: 94 additions & 0 deletions pkg/session/tidb_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,14 +22,108 @@ import (
"github.com/pingcap/tidb/pkg/kv"
"github.com/pingcap/tidb/pkg/meta"
"github.com/pingcap/tidb/pkg/parser/ast"
session_metrics "github.com/pingcap/tidb/pkg/session/metrics"
"github.com/pingcap/tidb/pkg/sessionctx/vardef"
"github.com/pingcap/tidb/pkg/sessiontxn"
"github.com/pingcap/tidb/pkg/store/mockstore"
"github.com/pingcap/tidb/pkg/util"
"github.com/pingcap/tidb/pkg/util/dbterror/exeerrors"
"github.com/pingcap/tidb/pkg/util/execdetails"
"github.com/pingcap/tidb/pkg/util/sqlexec"
"github.com/stretchr/testify/require"
)

type recordingObserver struct {
count int
}

func (o *recordingObserver) Observe(float64) {
o.count++
}

func TestSharedLockLostRollsBackTransaction(t *testing.T) {
store, dom := CreateStoreAndBootstrap(t)
defer func() { require.NoError(t, store.Close()) }()
defer dom.Close()

testCases := []struct {
name string
beginSQL string
pessimistic bool
sharedLockLost bool
}{
{
name: "pessimistic shared lock lost",
beginSQL: "begin pessimistic",
pessimistic: true,
sharedLockLost: true,
},
{
name: "optimistic shared lock lost mode mismatch",
beginSQL: "begin optimistic",
pessimistic: false,
sharedLockLost: true,
},
{
name: "pessimistic deadlock unchanged",
beginSQL: "begin pessimistic",
pessimistic: true,
sharedLockLost: false,
},
}

for _, testCase := range testCases {
t.Run(testCase.name, func(t *testing.T) {
se, err := createSession(store)
require.NoError(t, err)
defer se.Close()

MustExec(t, se, testCase.beginSQL)
txnManager := sessiontxn.GetTxnManager(se)
txn, err := txnManager.ActivateTxn()
require.NoError(t, err)
require.True(t, txn.Valid())
require.Equal(t, testCase.pessimistic, txn.IsPessimistic())
require.NotNil(t, txnManager.GetContextProvider())

recorder := &recordingObserver{}
var restoreObserver func()
if testCase.pessimistic {
originalObserver := session_metrics.TransactionDurationPessimisticAbortGeneral
session_metrics.TransactionDurationPessimisticAbortGeneral = recorder
restoreObserver = func() {
session_metrics.TransactionDurationPessimisticAbortGeneral = originalObserver
}
} else {
originalObserver := session_metrics.TransactionDurationOptimisticAbortGeneral
session_metrics.TransactionDurationOptimisticAbortGeneral = recorder
restoreObserver = func() {
session_metrics.TransactionDurationOptimisticAbortGeneral = originalObserver
}
}
defer restoreObserver()

var stmtErr error
if testCase.sharedLockLost {
stmtErr = kv.ErrSharedLockLost.GenWithStackByArgs(txn.StartTS(), "6B6579")
} else {
stmtErr = exeerrors.ErrDeadlock
}
got := autoCommitAfterStmt(context.Background(), se, stmtErr, nil)
require.Same(t, stmtErr, got)
if testCase.sharedLockLost {
require.True(t, kv.ErrSharedLockLost.Equal(got))
} else {
require.True(t, exeerrors.ErrDeadlock.Equal(got))
}
require.False(t, se.sessionVars.InTxn())
require.False(t, se.txn.Valid())
require.Nil(t, txnManager.GetContextProvider())
require.Equal(t, 1, recorder.count)
})
}
}

func TestDomapHandleNil(t *testing.T) {
// this is required for enterprise plugins
// ref: https://github.com/pingcap/tidb/issues/37319
Expand Down
Loading
Loading