diff --git a/error/error.go b/error/error.go index 792047b890..8060e5c437 100644 --- a/error/error.go +++ b/error/error.go @@ -134,6 +134,25 @@ func (d *ErrDeadlock) Error() string { return d.String() } +// ErrLockUpgradeConflict wraps *kvrpcpb.LockUpgradeConflict to implement the error interface. +type ErrLockUpgradeConflict struct { + *kvrpcpb.LockUpgradeConflict +} + +func (e *ErrLockUpgradeConflict) Error() string { + return fmt.Sprintf("lock upgrade conflict { %s }", e.String()) +} + +// ErrSharedLockLost wraps *kvrpcpb.SharedLockLost to report that a transaction +// deterministically lost its shared-lock ownership during an upgrade. +type ErrSharedLockLost struct { + *kvrpcpb.SharedLockLost +} + +func (e *ErrSharedLockLost) Error() string { + return fmt.Sprintf("shared lock lost { %s }", e.String()) +} + // PDError wraps *pdpb.Error to implement the error interface. type PDError struct { Err *pdpb.Error @@ -338,10 +357,18 @@ func ExtractKeyErr(keyErr *kvrpcpb.KeyError) error { redact.RedactKeyErrIfNecessary(keyErr) + if keyErr.SharedLockLost != nil { + return errors.WithStack(&ErrSharedLockLost{SharedLockLost: keyErr.SharedLockLost}) + } + if keyErr.Conflict != nil { return errors.WithStack(NewErrWriteConflict(keyErr.GetConflict())) } + if keyErr.LockUpgradeConflict != nil { + return errors.WithStack(&ErrLockUpgradeConflict{LockUpgradeConflict: keyErr.LockUpgradeConflict}) + } + if keyErr.Retryable != "" { return errors.WithStack(&ErrRetryable{Retryable: keyErr.Retryable}) } diff --git a/error/error_test.go b/error/error_test.go index 67eea837c8..d4f7e3d49b 100644 --- a/error/error_test.go +++ b/error/error_test.go @@ -1,13 +1,68 @@ package error import ( + stderrs "errors" "testing" "github.com/pingcap/errors" "github.com/pingcap/kvproto/pkg/kvrpcpb" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" ) +func TestExtractKeyErrLockUpgradeConflict(t *testing.T) { + keyErr := &kvrpcpb.KeyError{ + LockUpgradeConflict: &kvrpcpb.LockUpgradeConflict{ + Key: []byte("key"), + StartTs: 101, + OwnerStartTs: 202, + Reason: kvrpcpb.LockUpgradeConflict_SecondUpgrader, + }, + } + + err := ExtractKeyErr(keyErr) + require.Error(t, err) + require.False(t, IsErrWriteConflict(err)) + + var retryable *ErrRetryable + require.False(t, stderrs.As(err, &retryable)) + + var conflict *ErrLockUpgradeConflict + require.ErrorAs(t, err, &conflict) + require.Equal(t, []byte("key"), conflict.Key) + require.Equal(t, uint64(101), conflict.StartTs) + require.Equal(t, uint64(202), conflict.OwnerStartTs) + require.Equal(t, kvrpcpb.LockUpgradeConflict_SecondUpgrader, conflict.Reason) +} + +func TestExtractKeyErrSharedLockLost(t *testing.T) { + keyErr := &kvrpcpb.KeyError{ + SharedLockLost: &kvrpcpb.SharedLockLost{ + Key: []byte("key"), + StartTs: 101, + }, + Conflict: &kvrpcpb.WriteConflict{ + StartTs: 101, + Reason: kvrpcpb.WriteConflict_PessimisticRetry, + }, + } + + err := ExtractKeyErr(keyErr) + require.Error(t, err) + require.False(t, IsErrWriteConflict(err)) + require.False(t, IsErrorUndetermined(err)) + + var retryable *ErrRetryable + require.False(t, stderrs.As(err, &retryable)) + var deadlock *ErrDeadlock + require.False(t, stderrs.As(err, &deadlock)) + + var lost *ErrSharedLockLost + require.ErrorAs(t, err, &lost) + require.Equal(t, []byte("key"), lost.Key) + require.Equal(t, uint64(101), lost.StartTs) +} + func TestExtractDebugInfoStrFromKeyErr(t *testing.T) { origRedact := errors.RedactLogEnabled.Load() defer errors.RedactLogEnabled.Store(origRedact) diff --git a/go.mod b/go.mod index a62644e622..da1cf5bbe5 100644 --- a/go.mod +++ b/go.mod @@ -15,7 +15,7 @@ require ( github.com/pingcap/errors v0.11.5-0.20241219054535-6b8c588c3122 github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 github.com/pingcap/goleveldb v0.0.0-20191226122134-f82aafb29989 - 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.20221110025148-ca232912c9f3 github.com/pkg/errors v0.9.1 github.com/prometheus/client_golang v1.20.5 @@ -35,6 +35,8 @@ require ( modernc.org/mathutil v1.7.1 ) +replace github.com/pingcap/kvproto => github.com/wfxr/kvproto v0.0.0-20260806092442-d04fa0402753 + require ( github.com/beorn7/perks v1.0.1 // indirect github.com/cespare/xxhash/v2 v2.3.0 // indirect diff --git a/go.sum b/go.sum index 3bde9f254d..015521dfd9 100644 --- a/go.sum +++ b/go.sum @@ -81,8 +81,6 @@ github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 h1:tdMsjOqUR7YXH github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86/go.mod h1:exzhVYca3WRtd6gclGNErRWb1qEgff3LYta0LvRmON4= github.com/pingcap/goleveldb v0.0.0-20191226122134-f82aafb29989 h1:surzm05a8C9dN8dIUmo4Be2+pMRb6f55i+UIYrluu2E= github.com/pingcap/goleveldb v0.0.0-20191226122134-f82aafb29989/go.mod h1:O17XtbryoCJhkKGbT62+L2OlrniwqiGLSqrmdHCMzZw= -github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 h1:qW0gHsqY3X3qmyiEfESqpYSF3Vu7agUiAyBnPWQXtm8= -github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE= github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3 h1:HR/ylkkLmGdSSDaD8IDP+SZrdhV1Kibl9KrHxJ9eciw= github.com/pingcap/log v1.1.1-0.20221110025148-ca232912c9f3/go.mod h1:DWQW5jICDR7UJh4HtxXSM20Churx4CQL0fwL/SoOSA4= github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= @@ -119,6 +117,8 @@ github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3 h1:OoBvgoeWmdNEXtS+ github.com/tikv/pd/client v0.0.0-20260708075407-4e05b9d2c2d3/go.mod h1:3/Bu91CJONgkDA+Y0v/cnbROSJnu5tQ09vv7JGybUBA= github.com/twmb/murmur3 v1.1.3 h1:D83U0XYKcHRYwYIpBKf3Pks91Z0Byda/9SJ8B6EMRcA= github.com/twmb/murmur3 v1.1.3/go.mod h1:Qq/R7NUyOfr65zD+6Q5IHKsJLwP7exErjN6lyyq3OSQ= +github.com/wfxr/kvproto v0.0.0-20260806092442-d04fa0402753 h1:wQuil8SCJhSp+LJqqiMgPjzBZgzobJCIQFp52UIYq6Y= +github.com/wfxr/kvproto v0.0.0-20260806092442-d04fa0402753/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE= github.com/yuin/goldmark v1.1.27/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= github.com/yuin/goldmark v1.2.1/go.mod h1:3hX8gzYuyVAZsxl0MRgGTJEmQBFcNTphYh9decYSb74= go.etcd.io/etcd/api/v3 v3.5.10 h1:szRajuUUbLyppkhs9K6BRtjY37l66XQQmw7oZRANE4k= diff --git a/integration_tests/go.mod b/integration_tests/go.mod index 783e322d87..c615efb54a 100644 --- a/integration_tests/go.mod +++ b/integration_tests/go.mod @@ -7,7 +7,7 @@ require ( github.com/ninedraft/israce v0.0.3 github.com/pingcap/errors v0.11.5-0.20260508054701-306e305bcf41 github.com/pingcap/failpoint v0.0.0-20240528011301-b51a646c7c86 - github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 + github.com/pingcap/kvproto v0.0.0-20260806092442-d04fa0402753 github.com/pingcap/tidb v1.1.0-beta.0.20260715060322-10292a4f8697 github.com/pkg/errors v0.9.1 github.com/prometheus/client_golang v1.23.0 @@ -177,5 +177,6 @@ require ( replace ( 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/tikv/client-go/v2 => ../ ) diff --git a/integration_tests/go.sum b/integration_tests/go.sum index f0e172e567..654fd2bb8b 100644 --- a/integration_tests/go.sum +++ b/integration_tests/go.sum @@ -1451,9 +1451,6 @@ github.com/pingcap/fn v1.0.0 h1:CyA6AxcOZkQh52wIqYlAmaVmF6EvrcqFywP463pjA8g= github.com/pingcap/fn v1.0.0/go.mod h1:u9WZ1ZiOD1RpNhcI42RucFh/lBuzTu6rw88a+oF2Z24= github.com/pingcap/goleveldb v0.0.0-20191226122134-f82aafb29989 h1:surzm05a8C9dN8dIUmo4Be2+pMRb6f55i+UIYrluu2E= github.com/pingcap/goleveldb v0.0.0-20191226122134-f82aafb29989/go.mod h1:O17XtbryoCJhkKGbT62+L2OlrniwqiGLSqrmdHCMzZw= -github.com/pingcap/kvproto v0.0.0-20241113043844-e1fa7ea8c302/go.mod h1:rXxWk2UnwfUhLXha1jxRWPADw9eMZGWEWCg92Tgmb/8= -github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368 h1:qW0gHsqY3X3qmyiEfESqpYSF3Vu7agUiAyBnPWQXtm8= -github.com/pingcap/kvproto v0.0.0-20260721064811-683dad8fa368/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE= github.com/pingcap/log v0.0.0-20210625125904-98ed8e2eb1c7/go.mod h1:8AanEdAHATuRurdGxZXBz0At+9avep+ub7U1AGYLIMM= github.com/pingcap/log v1.1.0/go.mod h1:DWQW5jICDR7UJh4HtxXSM20Churx4CQL0fwL/SoOSA4= github.com/pingcap/log v1.1.1-0.20250917021125-19901e015dc9 h1:qG9BSvlWFEE5otQGamuWedx9LRm0nrHvsQRQiW8SxEs= @@ -1621,6 +1618,8 @@ github.com/vbauerster/mpb/v7 v7.5.3 h1:BkGfmb6nMrrBQDFECR/Q7RkKCw7ylMetCb4079CGs github.com/vbauerster/mpb/v7 v7.5.3/go.mod h1:i+h4QY6lmLvBNK2ah1fSreiw3ajskRlBp9AhY/PnuOE= github.com/wangjohn/quickselect v0.0.0-20161129230411-ed8402a42d5f h1:9DDCDwOyEy/gId+IEMrFHLuQ5R/WV0KNxWLler8X2OY= github.com/wangjohn/quickselect v0.0.0-20161129230411-ed8402a42d5f/go.mod h1:8sdOQnirw1PrcnTJYkmW1iOHtUmblMmGdUOHyWYycLI= +github.com/wfxr/kvproto v0.0.0-20260806092442-d04fa0402753 h1:wQuil8SCJhSp+LJqqiMgPjzBZgzobJCIQFp52UIYq6Y= +github.com/wfxr/kvproto v0.0.0-20260806092442-d04fa0402753/go.mod h1:z6+aAHB7dBkA+LyinEX+48/ImRJ3jag0Hg0c7wkhEvE= github.com/xiang90/probing v0.0.0-20221125231312-a49e3df8f510 h1:S2dVYn90KE98chqDkyE9Z4N61UnQd+KOfgp5Iu53llk= github.com/xiang90/probing v0.0.0-20221125231312-a49e3df8f510/go.mod h1:UETIi67q53MR2AWcXfiuqkDkRtnGDLqkBTpCHuJHxtU= github.com/xordataexchange/crypt v0.0.3-0.20170626215501-b2862e3d0a77/go.mod h1:aYKd//L2LvnjZzWKhF00oedf4jCCReLcmhLdhm1A27Q= @@ -1742,6 +1741,16 @@ golang.org/x/crypto v0.18.0/go.mod h1:R0j02AL6hcrfOiy9T4ZYp/rcWeMxM3L6QYxlOuEG1m golang.org/x/crypto v0.19.0/go.mod h1:Iy9bg/ha4yyC70EfRS8jz+B6ybOBKMaSxLj6P6oBDfU= golang.org/x/crypto v0.23.0/go.mod h1:CKFgDieR+mRhux2Lsu27y0fO304Db0wZe70UKqHu0v8= golang.org/x/crypto v0.24.0/go.mod h1:Z1PMYSOR5nyMcyAVAIQSKCDwalqy85Aqn1x3Ws4L5DM= +golang.org/x/crypto v0.38.0/go.mod h1:MvrbAqul58NNYPKnOra203SB9vpuZW0e+RRZV+Ggqjw= +golang.org/x/crypto v0.39.0/go.mod h1:L+Xg3Wf6HoL4Bn4238Z6ft6KfEpN0tJGo53AAPC632U= +golang.org/x/crypto v0.40.0/go.mod h1:Qr1vMER5WyS2dfPHAlsOj01wgLbsyWtFn/aY+5+ZdxY= +golang.org/x/crypto v0.41.0/go.mod h1:pO5AFd7FA68rFak7rOAGVuygIISepHftHnr8dr6+sUc= +golang.org/x/crypto v0.42.0/go.mod h1:4+rDnOTJhQCx2q7/j6rAN5XDw8kPjeaXEUR2eL94ix8= +golang.org/x/crypto v0.43.0/go.mod h1:BFbav4mRNlXJL4wNeejLpWxB7wMbc79PdRGhWKncxR0= +golang.org/x/crypto v0.44.0/go.mod h1:013i+Nw79BMiQiMsOPcVCB5ZIJbYkerPrGnOa00tvmc= +golang.org/x/crypto v0.46.0/go.mod h1:Evb/oLKmMraqjZ2iQTwDwvCtJkczlDuTmdJXoZVzqU0= +golang.org/x/crypto v0.47.0/go.mod h1:ff3Y9VzzKbwSSEzWqJsJVBnWmRwRSHt/6Op5n9bQc4A= +golang.org/x/crypto v0.48.0/go.mod h1:r0kV5h3qnFPlQnBSrULhlsRfryS2pmewsg+XfMgkVos= golang.org/x/crypto v0.51.0 h1:IBPXwPfKxY7cWQZ38ZCIRPI50YLeevDLlLnyC5wRGTI= golang.org/x/crypto v0.51.0/go.mod h1:8AdwkbraGNABw2kOX6YFPs3WM22XqI4EXEd8g+x7Oc8= golang.org/x/exp v0.0.0-20180321215751-8460e604b9de/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA= @@ -1808,6 +1817,15 @@ golang.org/x/mod v0.11.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= golang.org/x/mod v0.12.0/go.mod h1:iBbtSCu2XBx23ZKBPSOrRkjjQPZFPuis4dIYUhu/chs= golang.org/x/mod v0.15.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c= golang.org/x/mod v0.17.0/go.mod h1:hTbmBsO62+eylJbnUtE2MGJUyE7QWk4xUqPFrRgJ+7c= +golang.org/x/mod v0.24.0/go.mod h1:IXM97Txy2VM4PJ3gI61r1YEk/gAj6zAHN3AdZt6S9Ww= +golang.org/x/mod v0.25.0/go.mod h1:IXM97Txy2VM4PJ3gI61r1YEk/gAj6zAHN3AdZt6S9Ww= +golang.org/x/mod v0.26.0/go.mod h1:/j6NAhSk8iQ723BGAUyoAcn7SlD7s15Dp9Nd/SfeaFQ= +golang.org/x/mod v0.27.0/go.mod h1:rWI627Fq0DEoudcK+MBkNkCe0EetEaDSwJJkCcjpazc= +golang.org/x/mod v0.28.0/go.mod h1:yfB/L0NOf/kmEbXjzCPOx1iK1fRutOydrCMsqRhEBxI= +golang.org/x/mod v0.29.0/go.mod h1:NyhrlYXJ2H4eJiRy/WDBO6HMqZQ6q9nk4JzS3NuCK+w= +golang.org/x/mod v0.30.0/go.mod h1:lAsf5O2EvJeSFMiBxXDki7sCgAxEUcZHXoXMKT4GJKc= +golang.org/x/mod v0.31.0/go.mod h1:43JraMp9cGx1Rx3AqioxrbrhNsLl2l/iNAvuBkrezpg= +golang.org/x/mod v0.32.0/go.mod h1:SgipZ/3h2Ci89DlEtEXWUk/HteuRin+HHhN+WbNhguU= golang.org/x/mod v0.35.0 h1:Ww1D637e6Pg+Zb2KrWfHQUnH2dQRLBQyAtpr/haaJeM= golang.org/x/mod v0.35.0/go.mod h1:+GwiRhIInF8wPm+4AoT6L0FA1QWAad3OMdTRx4tFYlU= golang.org/x/net v0.0.0-20180724234803-3673e40ba225/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= @@ -1880,6 +1898,17 @@ golang.org/x/net v0.20.0/go.mod h1:z8BVo6PvndSri0LbOE3hAn0apkU+1YvI6E70E9jsnvY= golang.org/x/net v0.21.0/go.mod h1:bIjVDfnllIU7BJ2DNgfnXvpSvtn8VRwhlsaeUTyUS44= golang.org/x/net v0.25.0/go.mod h1:JkAGAh7GEvH74S6FOH42FLoXpXbE/aqXSrIQjXgsiwM= golang.org/x/net v0.26.0/go.mod h1:5YKkiSynbBIh3p6iOc/vibscux0x38BZDkn8sCUPxHE= +golang.org/x/net v0.40.0/go.mod h1:y0hY0exeL2Pku80/zKK7tpntoX23cqL3Oa6njdgRtds= +golang.org/x/net v0.41.0/go.mod h1:B/K4NNqkfmg07DQYrbwvSluqCJOOXwUjeb/5lOisjbA= +golang.org/x/net v0.42.0/go.mod h1:FF1RA5d3u7nAYA4z2TkclSCKh68eSXtiFwcWQpPXdt8= +golang.org/x/net v0.43.0/go.mod h1:vhO1fvI4dGsIjh73sWfUVjj3N7CA9WkKJNQm2svM6Jg= +golang.org/x/net v0.44.0/go.mod h1:ECOoLqd5U3Lhyeyo/QDCEVQ4sNgYsqvCZ722XogGieY= +golang.org/x/net v0.45.0/go.mod h1:ECOoLqd5U3Lhyeyo/QDCEVQ4sNgYsqvCZ722XogGieY= +golang.org/x/net v0.46.0/go.mod h1:Q9BGdFy1y4nkUwiLvT5qtyhAnEHgnQ/zd8PfU6nc210= +golang.org/x/net v0.47.0/go.mod h1:/jNxtkgq5yWUGYkaZGqo27cfGZ1c5Nen03aYrrKpVRU= +golang.org/x/net v0.48.0/go.mod h1:+ndRgGjkh8FGtu1w1FGbEC31if4VrNVMuKTgcAAnQRY= +golang.org/x/net v0.49.0/go.mod h1:/ysNB2EvaqvesRkuLAyjI1ycPZlQHM3q01F02UY/MV8= +golang.org/x/net v0.51.0/go.mod h1:aamm+2QF5ogm02fjy5Bb7CQ0WMt1/WVM7FtyaTLlA9Y= golang.org/x/net v0.54.0 h1:2zJIZAxAHV/OHCDTCOHAYehQzLfSXuf/5SoL/Dv6w/w= golang.org/x/net v0.54.0/go.mod h1:Sj4oj8jK6XmHpBZU/zWHw3BV3abl4Kvi+Ut7cQcY+cQ= golang.org/x/oauth2 v0.0.0-20180821212333-d2e6202438be/go.mod h1:N/0e6XlmueqKjAGxoOufVs8QHGRruUQn6yWY3a++T0U= @@ -1936,6 +1965,12 @@ golang.org/x/sync v0.2.0/go.mod h1:RxMgew5VJxzue5/jJTE5uejpjVlOe/izrB70Jof72aM= golang.org/x/sync v0.3.0/go.mod h1:FU7BRWz2tNW+3quACPkgCx/L+uEAv1htQ0V83Z9Rj+Y= golang.org/x/sync v0.6.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= golang.org/x/sync v0.7.0/go.mod h1:Czt+wKu1gCyEFDUtn0jG5QVvpJ6rzVqr5aXyt9drQfk= +golang.org/x/sync v0.14.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= +golang.org/x/sync v0.15.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= +golang.org/x/sync v0.16.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= +golang.org/x/sync v0.17.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +golang.org/x/sync v0.18.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= +golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI= golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4= golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.0.0-20180830151530-49385e6e1522/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= @@ -2034,9 +2069,26 @@ golang.org/x/sys v0.16.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.17.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.20.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= golang.org/x/sys v0.21.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA= +golang.org/x/sys v0.33.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= +golang.org/x/sys v0.34.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= +golang.org/x/sys v0.35.0/go.mod h1:BJP2sWEmIv4KK5OTEluFJCKSidICx8ciO85XgH3Ak8k= +golang.org/x/sys v0.36.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/sys v0.37.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/sys v0.38.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/sys v0.39.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/sys v0.40.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= +golang.org/x/sys v0.41.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks= golang.org/x/sys v0.44.0 h1:ildZl3J4uzeKP07r2F++Op7E9B29JRUy+a27EibtBTQ= golang.org/x/sys v0.44.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/telemetry v0.0.0-20240228155512-f48c80bd79b2/go.mod h1:TeRTkGYfJXctD9OcfyVLyj2J3IxLnKwHJR8f4D8a3YE= +golang.org/x/telemetry v0.0.0-20240521205824-bda55230c457/go.mod h1:pRgIJT+bRLFKnoM1ldnzKoxTIn14Yxz928LQRYYgIN0= +golang.org/x/telemetry v0.0.0-20250710130107-8d8967aff50b/go.mod h1:4ZwOYna0/zsOKwuR5X/m0QFOJpSZvAxFfkQT+Erd9D4= +golang.org/x/telemetry v0.0.0-20250807160809-1a19826ec488/go.mod h1:fGb/2+tgXXjhjHsTNdVEEMZNWA0quBnfrO+AfoDSAKw= +golang.org/x/telemetry v0.0.0-20250908211612-aef8a434d053/go.mod h1:+nZKN+XVh4LCiA9DV3ywrzN4gumyCnKjau3NGb9SGoE= +golang.org/x/telemetry v0.0.0-20251008203120-078029d740a8/go.mod h1:Pi4ztBfryZoJEkyFTI5/Ocsu2jXyDr6iSdgJiYE/uwE= +golang.org/x/telemetry v0.0.0-20251111182119-bc8e575c7b54/go.mod h1:hKdjCMrbv9skySur+Nek8Hd0uJ0GuxJIoIX2payrIdQ= +golang.org/x/telemetry v0.0.0-20251203150158-8fff8a5912fc/go.mod h1:hKdjCMrbv9skySur+Nek8Hd0uJ0GuxJIoIX2payrIdQ= +golang.org/x/telemetry v0.0.0-20260109210033-bd525da824e2/go.mod h1:b7fPSJ0pKZ3ccUh8gnTONJxhn3c/PS6tyzQvyqw4iA8= golang.org/x/telemetry v0.0.0-20260409153401-be6f6cb8b1fa h1:efT73AJZfAAUV7SOip6pWGkwJDzIGiKBZGVzHYa+ve4= golang.org/x/telemetry v0.0.0-20260409153401-be6f6cb8b1fa/go.mod h1:kHjTxDEnAu6/Nl9lDkzjWpR+bmKfxeiRuSDlsMb70gE= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= @@ -2058,6 +2110,15 @@ golang.org/x/term v0.16.0/go.mod h1:yn7UURbUtPyrVJPGPq404EukNFxcm/foM+bV/bfcDsY= golang.org/x/term v0.17.0/go.mod h1:lLRBjIVuehSbZlaOtGMbcMncT+aqLLLmKrsjNrUguwk= golang.org/x/term v0.20.0/go.mod h1:8UkIAJTvZgivsXaD6/pH6U9ecQzZ45awqEOzuCvwpFY= golang.org/x/term v0.21.0/go.mod h1:ooXLefLobQVslOqselCNF4SxFAaoS6KujMbsGzSDmX0= +golang.org/x/term v0.32.0/go.mod h1:uZG1FhGx848Sqfsq4/DlJr3xGGsYMu/L5GW4abiaEPQ= +golang.org/x/term v0.33.0/go.mod h1:s18+ql9tYWp1IfpV9DmCtQDDSRBUjKaw9M1eAv5UeF0= +golang.org/x/term v0.34.0/go.mod h1:5jC53AEywhIVebHgPVeg0mj8OD3VO9OzclacVrqpaAw= +golang.org/x/term v0.35.0/go.mod h1:TPGtkTLesOwf2DE8CgVYiZinHAOuy5AYUYT1lENIZnA= +golang.org/x/term v0.36.0/go.mod h1:Qu394IJq6V6dCBRgwqshf3mPF85AqzYEzofzRdZkWss= +golang.org/x/term v0.37.0/go.mod h1:5pB4lxRNYYVZuTLmy8oR2BH8dflOR+IbTYFD8fi3254= +golang.org/x/term v0.38.0/go.mod h1:bSEAKrOT1W+VSu9TSCMtoGEOUcKxOKgl3LE5QEF/xVg= +golang.org/x/term v0.39.0/go.mod h1:yxzUCTP/U+FzoxfdKmLaA0RV1WgE0VY7hXBwKtY/4ww= +golang.org/x/term v0.40.0/go.mod h1:w2P8uVp06p2iyKKuvXIm7N/y0UCRt3UfJTfZ7oOpglM= golang.org/x/term v0.43.0 h1:S4RLU2sB31O/NCl+zFN9Aru9A/Cq2aqKpTZJ6B+DwT4= golang.org/x/term v0.43.0/go.mod h1:lrhlHNdQJHO+1qVYiHfFKVuVioJIheAc3fBSMFYEIsk= golang.org/x/text v0.0.0-20170915032832-14c0d48ead0c/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ= @@ -2083,6 +2144,16 @@ golang.org/x/text v0.13.0/go.mod h1:TvPlkZtksWOMsz7fbANvkp4WM8x/WCo/om8BMLbz+aE= golang.org/x/text v0.14.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= golang.org/x/text v0.15.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU= golang.org/x/text v0.16.0/go.mod h1:GhwF1Be+LQoKShO3cGOHzqOgRrGaYc9AvblQOmPVHnI= +golang.org/x/text v0.25.0/go.mod h1:WEdwpYrmk1qmdHvhkSTNPm3app7v4rsT8F2UD6+VHIA= +golang.org/x/text v0.26.0/go.mod h1:QK15LZJUUQVJxhz7wXgxSy/CJaTFjd0G+YLonydOVQA= +golang.org/x/text v0.27.0/go.mod h1:1D28KMCvyooCX9hBiosv5Tz/+YLxj0j7XhWjpSUF7CU= +golang.org/x/text v0.28.0/go.mod h1:U8nCwOR8jO/marOQ0QbDiOngZVEBB7MAiitBuMjXiNU= +golang.org/x/text v0.29.0/go.mod h1:7MhJOA9CD2qZyOKYazxdYMF85OwPdEr9jTtBpO7ydH4= +golang.org/x/text v0.30.0/go.mod h1:yDdHFIX9t+tORqspjENWgzaCVXgk0yYnYuSZ8UzzBVM= +golang.org/x/text v0.31.0/go.mod h1:tKRAlv61yKIjGGHX/4tP1LTbc13YSec1pxVEWXzfoeM= +golang.org/x/text v0.32.0/go.mod h1:o/rUWzghvpD5TXrTIBuJU77MTaN0ljMWE47kxGJQ7jY= +golang.org/x/text v0.33.0/go.mod h1:LuMebE6+rBincTi9+xWTY8TztLzKHc/9C1uBCG27+q8= +golang.org/x/text v0.34.0/go.mod h1:homfLqTYRFyVYemLBFl5GgL/DWEiH5wcsQ5gSh1yziA= golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc= golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38= golang.org/x/time v0.0.0-20181108054448-85acf8d2951c/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ= @@ -2166,6 +2237,15 @@ golang.org/x/tools v0.9.1/go.mod h1:owI94Op576fPu3cIGQeHs3joujW/2Oc6MtlxbF5dfNc= golang.org/x/tools v0.10.0/go.mod h1:UJwyiVBsOA2uwvK/e5OY3GTpDUJriEd+/YlqAwLPmyM= golang.org/x/tools v0.13.0/go.mod h1:HvlwmtVNQAhOuCjW7xxvovg8wbNq7LwfXh/k7wXUl58= golang.org/x/tools v0.21.1-0.20240508182429-e35e4ccd0d2d/go.mod h1:aiJjzUbINMkxbQROHiO6hDPo2LHcIPhhQsa9DLh0yGk= +golang.org/x/tools v0.33.0/go.mod h1:CIJMaWEY88juyUfo7UbgPqbC8rU2OqfAV1h2Qp0oMYI= +golang.org/x/tools v0.34.0/go.mod h1:pAP9OwEaY1CAW3HOmg3hLZC5Z0CCmzjAF2UQMSqNARg= +golang.org/x/tools v0.35.0/go.mod h1:NKdj5HkL/73byiZSJjqJgKn3ep7KjFkBOkR/Hps3VPw= +golang.org/x/tools v0.36.0/go.mod h1:WBDiHKJK8YgLHlcQPYQzNCkUxUypCaa5ZegCVutKm+s= +golang.org/x/tools v0.37.0/go.mod h1:MBN5QPQtLMHVdvsbtarmTNukZDdgwdwlO5qGacAzF0w= +golang.org/x/tools v0.38.0/go.mod h1:yEsQ/d/YK8cjh0L6rZlY8tgtlKiBNTL14pGDJPJpYQs= +golang.org/x/tools v0.39.0/go.mod h1:JnefbkDPyD8UU2kI5fuf8ZX4/yUeh9W877ZeBONxUqQ= +golang.org/x/tools v0.40.0/go.mod h1:Ik/tzLRlbscWpqqMRjyWYDisX8bG13FrdXp3o4Sr9lc= +golang.org/x/tools v0.41.0/go.mod h1:XSY6eDqxVNiYgezAVqqCeihT4j1U2CCsqvH3WhQpnlg= golang.org/x/tools v0.44.0 h1:UP4ajHPIcuMjT1GqzDWRlalUEoY+uzoZKnhOjbIPD2c= golang.org/x/tools v0.44.0/go.mod h1:KA0AfVErSdxRZIsOVipbv3rQhVXTnlU6UhKxHd1seDI= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= diff --git a/internal/apicodec/codec_v2.go b/internal/apicodec/codec_v2.go index 070e75fe91..e24d42fc8c 100644 --- a/internal/apicodec/codec_v2.go +++ b/internal/apicodec/codec_v2.go @@ -1033,6 +1033,18 @@ func (c *codecV2) decodeKeyError(keyError *kvrpcpb.KeyError) (*kvrpcpb.KeyError, } } } + if keyError.SharedLockLost != nil { + keyError.SharedLockLost.Key, err = c.DecodeKey(keyError.SharedLockLost.Key) + if err != nil { + return nil, err + } + } + if keyError.LockUpgradeConflict != nil { + keyError.LockUpgradeConflict.Key, err = c.DecodeKey(keyError.LockUpgradeConflict.Key) + if err != nil { + return nil, err + } + } if keyError.CommitTsExpired != nil { keyError.CommitTsExpired.Key, err = c.DecodeKey(keyError.CommitTsExpired.Key) if err != nil { diff --git a/internal/apicodec/codec_v2_test.go b/internal/apicodec/codec_v2_test.go index f92ba59ee0..2b889c6b75 100644 --- a/internal/apicodec/codec_v2_test.go +++ b/internal/apicodec/codec_v2_test.go @@ -530,6 +530,36 @@ func (suite *testCodecV2Suite) TestDecodeKeyError() { }, decoded.DebugInfo.MvccInfo[0].Mvcc.Lock.Secondaries) }, }, + { + name: "SharedLockLost", + err: &kvrpcpb.KeyError{ + SharedLockLost: &kvrpcpb.SharedLockLost{ + Key: append(keyspacePrefix, []byte("key1")...), + StartTs: 11, + }, + }, + validate: func(re *require.Assertions, decoded *kvrpcpb.KeyError) { + re.Equal([]byte("key1"), decoded.SharedLockLost.Key) + re.Equal(uint64(11), decoded.SharedLockLost.StartTs) + }, + }, + { + name: "LockUpgradeConflict", + err: &kvrpcpb.KeyError{ + LockUpgradeConflict: &kvrpcpb.LockUpgradeConflict{ + Key: append(keyspacePrefix, []byte("key1")...), + StartTs: 11, + OwnerStartTs: 22, + Reason: kvrpcpb.LockUpgradeConflict_DuplicateInFlight, + }, + }, + validate: func(re *require.Assertions, decoded *kvrpcpb.KeyError) { + re.Equal([]byte("key1"), decoded.LockUpgradeConflict.Key) + re.Equal(uint64(11), decoded.LockUpgradeConflict.StartTs) + re.Equal(uint64(22), decoded.LockUpgradeConflict.OwnerStartTs) + re.Equal(kvrpcpb.LockUpgradeConflict_DuplicateInFlight, decoded.LockUpgradeConflict.Reason) + }, + }, } codec := suite.codec diff --git a/kv/kv.go b/kv/kv.go index 4bd086183e..0d07b75a7a 100644 --- a/kv/kv.go +++ b/kv/kv.go @@ -67,6 +67,7 @@ type LockCtx struct { CheckExistence bool LockOnlyIfExists bool InShareMode bool + AllowSharedLockUpgrade bool Values map[string]ReturnedValue MaxLockedWithConflictTS uint64 ValuesLock sync.Mutex diff --git a/txnkv/transaction/2pc.go b/txnkv/transaction/2pc.go index f92051ae31..3554db8711 100644 --- a/txnkv/transaction/2pc.go +++ b/txnkv/transaction/2pc.go @@ -148,6 +148,7 @@ type twoPhaseCommitter struct { mu struct { sync.RWMutex undeterminedErr error // undeterminedErr saves the rpc error we encounter when commit primary key. + fatalTxnErr error committed bool } syncLog bool @@ -2424,6 +2425,23 @@ func (c *twoPhaseCommitter) getUndeterminedErr() error { return c.mu.undeterminedErr } +func (c *twoPhaseCommitter) setFatalTxnErr(err error) { + if err == nil { + return + } + c.mu.Lock() + defer c.mu.Unlock() + if c.mu.fatalTxnErr == nil { + c.mu.fatalTxnErr = err + } +} + +func (c *twoPhaseCommitter) getTxnStateErrs() (error, error) { + c.mu.RLock() + defer c.mu.RUnlock() + return c.mu.undeterminedErr, c.mu.fatalTxnErr +} + func (c *twoPhaseCommitter) mutationsOfKeys(keys [][]byte) CommitterMutations { var res PlainMutations for i := 0; i < c.mutations.Len(); i++ { diff --git a/txnkv/transaction/txn.go b/txnkv/transaction/txn.go index 52bb657fba..cd8c2d3184 100644 --- a/txnkv/transaction/txn.go +++ b/txnkv/transaction/txn.go @@ -799,6 +799,12 @@ func (txn *KVTxn) Commit(ctx context.Context) error { if !txn.valid { return tikverr.ErrInvalidTxn } + if txn.committer != nil { + undeterminedErr, fatalTxnErr := txn.committer.getTxnStateErrs() + if undeterminedErr == nil && fatalTxnErr != nil { + return fatalTxnErr + } + } defer txn.close() ctx = context.WithValue(ctx, util.RequestSourceKey, *txn.RequestSource) @@ -850,6 +856,9 @@ func (txn *KVTxn) Commit(ctx context.Context) error { } txn.committer = committer } + if err := txn.getTxnStateErr(); err != nil { + return err + } committer.SetDiskFullOpt(txn.diskFullOpt) committer.SetTxnSource(txn.txnSource) @@ -1384,6 +1393,224 @@ func (txn *KVTxn) LockKeysFunc(ctx context.Context, lockCtx *tikv.LockCtx, fn fu return txn.lockKeys(ctx, lockCtx, fn, keysInput...) } +func (txn *KVTxn) getTxnStateErr() error { + if txn.committer == nil { + return nil + } + undeterminedErr, fatalTxnErr := txn.committer.getTxnStateErrs() + if undeterminedErr != nil { + return errors.WithStack(tikverr.ErrResultUndetermined) + } + return fatalTxnErr +} + +// isLockUpgradeResultUndetermined reports whether a shared-to-exclusive +// upgrade error leaves the remote lock state uncertain. +// +// Upgrade requests are special because the transaction already owns the shared +// lock locally. Determined terminal errors are handled before this classifier. +// If TiKV returns a deterministic conflict or validation error, we know the new +// exclusive lock was not granted, so the caller should return that error +// directly and keep the original shared-lock state. Any error outside that +// known-safe set is treated as undetermined and must poison the transaction to +// avoid continuing with divergent local/remote lock state. +func isLockUpgradeResultUndetermined(err error) bool { + if err == nil { + return false + } + if tikverr.IsErrWriteConflict(err) || + tikverr.IsErrKeyExist(err) || + errors.Is(err, tikverr.ErrLockAcquireFailAndNoWaitSet) || + errors.Is(err, tikverr.ErrLockWaitTimeout) { + return false + } + var deadlock *tikverr.ErrDeadlock + if errors.As(err, &deadlock) { + return false + } + var lockUpgradeConflict *tikverr.ErrLockUpgradeConflict + if errors.As(err, &lockUpgradeConflict) { + return false + } + var assertionFailed *tikverr.ErrAssertionFailed + return !errors.As(err, &assertionFailed) +} + +func (txn *KVTxn) lockPessimisticKeyGroup( + ctx context.Context, + lockCtx *tikv.LockCtx, + keys [][]byte, + isUpgrade bool, +) (int, error) { + bo := retry.NewBackofferWithVars(ctx, pessimisticLockMaxBackoff, txn.vars) + txn.committer.isFirstLock = txn.lockedCnt == 0 && len(keys) == 1 + err := txn.committer.pessimisticLockMutations(bo, lockCtx, kvrpcpb.PessimisticLockWakeUpMode_WakeUpModeNormal, &PlainMutations{keys: keys}) + if lockCtx.Stats != nil && bo.GetTotalSleep() > 0 { + atomic.AddInt64(&lockCtx.Stats.BackoffTime, int64(bo.GetTotalSleep())*int64(time.Millisecond)) + lockCtx.Stats.Mu.Lock() + lockCtx.Stats.Mu.BackoffTypes = append(lockCtx.Stats.Mu.BackoffTypes, bo.GetTypes()...) + lockCtx.Stats.Mu.Unlock() + } + if err != nil { + var unmarkKeys [][]byte + memBuf := txn.us.GetMemBuffer() + memBuf.RLock() + for _, key := range keys { + if txn.us.HasPresumeKeyNotExists(key) { + unmarkKeys = append(unmarkKeys, key) + } + } + memBuf.RUnlock() + for _, key := range unmarkKeys { + txn.us.UnmarkPresumeKeyNotExists(key) + } + + if isUpgrade { + var sharedLockLost *tikverr.ErrSharedLockLost + if errors.As(err, &sharedLockLost) { + txn.committer.setFatalTxnErr(err) + return 0, err + } + if isLockUpgradeResultUndetermined(err) { + txn.committer.setUndeterminedErr(err) + return 0, errors.WithStack(tikverr.ErrResultUndetermined) + } + return 0, err + } + + keyMayBeLocked := !tikverr.IsErrWriteConflict(err) && !tikverr.IsErrKeyExist(err) + if len(keys) > 1 || keyMayBeLocked { + dl, isDeadlock := errors.Cause(err).(*tikverr.ErrDeadlock) + if isDeadlock { + if hashInKeys(dl.DeadlockKeyHash, keys) { + dl.IsRetryable = true + } + if lockCtx.OnDeadlock != nil { + lockCtx.OnDeadlock(dl) + } + } + + rollbackForUpdateTS := lockCtx.ForUpdateTS + if lockCtx.MaxLockedWithConflictTS > rollbackForUpdateTS { + rollbackForUpdateTS = lockCtx.MaxLockedWithConflictTS + } + wg := txn.asyncPessimisticRollback(ctx, keys, rollbackForUpdateTS) + + if isDeadlock { + logutil.Logger(ctx).Debug("deadlock error received", zap.Uint64("startTS", txn.startTS), zap.Stringer("deadlockInfo", dl)) + if dl.IsRetryable { + wg.Wait() + time.Sleep(time.Millisecond * 5) + if _, err := util.EvalFailpoint("SingleStmtDeadLockRetrySleep"); err == nil { + time.Sleep(300 * time.Millisecond) + } + } + } + } + return 0, err + } + + checkedExistence := lockCtx.CheckExistence + skippedLockKeys := 0 + memBuf := txn.us.GetMemBuffer() + for _, key := range keys { + valExists := true + if val, ok := lockCtx.Values[string(key)]; ok { + if lockCtx.ReturnValues || checkedExistence || val.LockedWithConflictTS != 0 { + if !val.Exists { + valExists = false + } + } + } + + if lockCtx.LockOnlyIfExists && !valExists { + skippedLockKeys++ + continue + } + + setValExists := tikv.SetKeyLockedValueExists + if !valExists { + setValExists = tikv.SetKeyLockedValueNotExists + } + memBuf.UpdateFlags(key, tikv.SetKeyLocked, tikv.DelNeedCheckExists, setValExists, tikv.SetKeyLockedInExclusiveMode) + } + if !isUpgrade { + txn.lockedCnt += len(keys) - skippedLockKeys + } + return skippedLockKeys, nil +} + +func (txn *KVTxn) lockKeysWithSharedLockUpgrade( + ctx context.Context, + lockCtx *tikv.LockCtx, + normalExclusiveKeys [][]byte, + upgradeKeys [][]byte, +) error { + if txn.committer == nil { + var sessionID uint64 + val := ctx.Value(util.SessionID) + if val != nil { + sessionID = val.(uint64) + } + var err error + txn.committer, err = newTwoPhaseCommitter(txn, sessionID) + if err != nil { + return err + } + } + + assignedPrimaryKey := false + totalKeys := len(normalExclusiveKeys) + len(upgradeKeys) + if txn.committer.primaryKey == nil { + // Prefer selecting the primary from freshly requested exclusive locks + // when possible, so the primary is not just a shared-locked key that + // still depends on upgrade success. + assignedPrimaryKey = true + keysForPrimary := normalExclusiveKeys + if len(keysForPrimary) == 0 { + keysForPrimary = upgradeKeys + } + txn.selectPrimaryForPessimisticLock(keysForPrimary) + } + + txn.committer.forUpdateTS = lockCtx.ForUpdateTS + lockCtx.Stats = &util.LockKeysDetails{ + LockKeys: int32(totalKeys), + ResolveLock: util.ResolveLockDetail{}, + } + + lockedInThisCall := 0 + if len(normalExclusiveKeys) > 0 { + skipped, err := txn.lockPessimisticKeyGroup(ctx, lockCtx, normalExclusiveKeys, false) + if err != nil { + if assignedPrimaryKey && lockedInThisCall == 0 { + txn.resetPrimary(false) + } + return err + } + lockedInThisCall += len(normalExclusiveKeys) - skipped + } + + for _, key := range upgradeKeys { + skipped, err := txn.lockPessimisticKeyGroup(ctx, lockCtx, [][]byte{key}, true) + if err != nil { + if assignedPrimaryKey && lockedInThisCall == 0 { + txn.resetPrimary(false) + } + return err + } + lockedInThisCall += 1 - skipped + } + + if assignedPrimaryKey && lockCtx.LockOnlyIfExists { + if totalKeys != 1 { + panic("LockOnlyIfExists only assigns the primary key when locking only one key") + } + txn.unsetPrimaryKeyIfNeeded(lockCtx) + } + return nil +} + func (txn *KVTxn) lockKeys(ctx context.Context, lockCtx *tikv.LockCtx, fn func(), keysInput ...[]byte) error { if txn.interceptor != nil { // User has called txn.SetRPCInterceptor() to explicitly set an interceptor, we @@ -1404,6 +1631,9 @@ func (txn *KVTxn) lockKeys(ctx context.Context, lockCtx *tikv.LockCtx, fn func() if err != nil { return err } + if err := txn.getTxnStateErr(); err != nil { + return err + } defer func() { if lockCtx.InShareMode { @@ -1456,6 +1686,7 @@ func (txn *KVTxn) lockKeys(ctx context.Context, lockCtx *tikv.LockCtx, fn func() memBuf := txn.us.GetMemBuffer() // Avoid data race with concurrent updates to the memBuf memBuf.RLock() + upgradeKeys := make([][]byte, 0, len(keysInput)) for _, key := range keysInput { // The value of lockedMap is only used by pessimistic transactions. var valueExist, locked, lockedInShareMode, checkKeyExists bool @@ -1466,9 +1697,16 @@ func (txn *KVTxn) lockKeys(ctx context.Context, lockCtx *tikv.LockCtx, fn func() checkKeyExists = flags.HasNeedCheckExists() } - if lockedInShareMode && !lockCtx.InShareMode { - memBuf.RUnlock() - return errors.New("upgrading a shared lock to an exclusive lock is not supported") + upgradeCandidate := lockedInShareMode && !lockCtx.InShareMode + if upgradeCandidate { + if !txn.IsPessimistic() || !lockCtx.AllowSharedLockUpgrade { + memBuf.RUnlock() + return errors.New("upgrading a shared lock to an exclusive lock is not supported") + } + if txn.IsInAggressiveLockingMode() { + memBuf.RUnlock() + return errors.New("shared lock upgrade is not supported in aggressive/fair locking mode") + } } // If the key is locked in the current aggressive locking stage, override the information in memBuf. @@ -1484,7 +1722,9 @@ func (txn *KVTxn) lockKeys(ctx context.Context, lockCtx *tikv.LockCtx, fn func() } } - if !locked || isInLastAggressiveLockingStage { + if upgradeCandidate { + upgradeKeys = append(upgradeKeys, key) + } else if !locked || isInLastAggressiveLockingStage { // Locks acquired in the previous aggressive locking stage might need to be updated later in // `filterAggressiveLockedKeys`. keys = append(keys, key) @@ -1497,7 +1737,7 @@ func (txn *KVTxn) lockKeys(ctx context.Context, lockCtx *tikv.LockCtx, fn func() return txn.committer.extractKeyExistsErr(e) } } - if lockCtx.ReturnValues && locked { + if lockCtx.ReturnValues && locked && !upgradeCandidate { keyStr := string(key) // An already locked key can not return values, we add an entry to let the caller get the value // in other ways. @@ -1506,28 +1746,41 @@ func (txn *KVTxn) lockKeys(ctx context.Context, lockCtx *tikv.LockCtx, fn func() } memBuf.RUnlock() - if len(keys) == 0 { + if len(keys) == 0 && len(upgradeKeys) == 0 { return nil } if lockCtx.LockOnlyIfExists { + var lockKey []byte + if len(keys) > 0 { + lockKey = keys[0] + } else { + lockKey = upgradeKeys[0] + } if !lockCtx.ReturnValues { return &tikverr.ErrLockOnlyIfExistsNoReturnValue{ StartTS: txn.startTS, ForUpdateTs: lockCtx.ForUpdateTS, - LockKey: keys[0], + LockKey: lockKey, } } // It can't transform LockOnlyIfExists mode to normal mode. If so, it can add a lock to a key // which doesn't exist in tikv. TiDB should ensure that primary key must be set when it sends // a LockOnlyIfExists pessimistic lock request. - if (txn.committer == nil || txn.committer.primaryKey == nil) && len(keys) > 1 { + if (txn.committer == nil || txn.committer.primaryKey == nil) && len(keys)+len(upgradeKeys) > 1 { return &tikverr.ErrLockOnlyIfExistsNoPrimaryKey{ StartTS: txn.startTS, ForUpdateTs: lockCtx.ForUpdateTS, - LockKey: keys[0], + LockKey: lockKey, } } } + if len(upgradeKeys) > 0 { + if len(keys) > 0 { + keys = deduplicateKeys(keys) + } + upgradeKeys = deduplicateKeys(upgradeKeys) + return txn.lockKeysWithSharedLockUpgrade(ctx, lockCtx, keys, upgradeKeys) + } keys = deduplicateKeys(keys) checkedExistence := false filteredAggressiveLockedKeysCount := 0 diff --git a/txnkv/transaction/txn_test.go b/txnkv/transaction/txn_test.go index b6d3b998ca..c4b929a1bc 100644 --- a/txnkv/transaction/txn_test.go +++ b/txnkv/transaction/txn_test.go @@ -16,6 +16,7 @@ package transaction import ( "context" + stderrs "errors" "testing" "time" @@ -23,6 +24,8 @@ import ( "github.com/pingcap/kvproto/pkg/metapb" "github.com/pkg/errors" "github.com/stretchr/testify/require" + "github.com/tikv/client-go/v2/config/retry" + tikverr "github.com/tikv/client-go/v2/error" "github.com/tikv/client-go/v2/internal/client" "github.com/tikv/client-go/v2/internal/locate" "github.com/tikv/client-go/v2/kv" @@ -88,6 +91,14 @@ func (m *mockStore) GetTiKVClient() client.Client { return &m.client } +type recorderStore struct { + *mockStore +} + +func (m *recorderStore) SendReq(bo *retry.Backoffer, req *tikvrpc.Request, regionID locate.RegionVerID, timeout time.Duration) (*tikvrpc.Response, error) { + return m.client.SendRequest(bo.GetCtx(), "mock-store", req, timeout) +} + func (m *mockStore) GetOracle() oracle.Oracle { return nil } @@ -273,8 +284,12 @@ func TestLockKeys(t *testing.T) { // `lockKeys` in exclusive mode again on k1 to upgrade the lock. lockCtx.InShareMode = false expectedLockType = kvrpcpb.Op_PessimisticLock - err = txn.lockKeys(context.TODO(), lockCtx, nil, key1) - require.ErrorContains(t, err, "upgrading a shared lock to an exclusive lock is not supported") + lockCtx.AllowSharedLockUpgrade = true + require.NoError(t, txn.lockKeys(context.TODO(), lockCtx, nil, key1)) + flags, err = txn.GetMemBuffer().GetFlags(key1) + require.NoError(t, err) + require.True(t, flags.HasLocked()) + require.False(t, flags.HasLockedInShareMode()) done := make(chan struct{}) go func() { @@ -289,6 +304,491 @@ func TestLockKeys(t *testing.T) { }) } +func TestSharedLockUpgrade(t *testing.T) { + type requestSummary struct { + cmd tikvrpc.CmdType + keys [][]byte + op kvrpcpb.Op + } + + newRecorderTxn := func(t *testing.T, onLock func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error)) (*testTxn, *[]requestSummary) { + txn := newTestTxn(t, 1) + txn.KVTxn.store = &recorderStore{mockStore: txn.store} + txn.SetPessimistic(true) + requests := make([]requestSummary, 0, 8) + txn.store.client.onSend = func(ctx context.Context, addr string, req *tikvrpc.Request, timeout time.Duration) (*tikvrpc.Response, error) { + switch req.Type { + case tikvrpc.CmdPessimisticLock: + lockReq := req.PessimisticLock() + keys := make([][]byte, len(lockReq.Mutations)) + for i, mutation := range lockReq.Mutations { + keys[i] = append([]byte(nil), mutation.Key...) + if i == 0 { + requests = append(requests, requestSummary{ + cmd: req.Type, + keys: keys[:0], + op: mutation.Op, + }) + } else { + require.Equal(t, requests[len(requests)-1].op, mutation.Op) + } + requests[len(requests)-1].keys = append(requests[len(requests)-1].keys, keys[i]) + } + return onLock(len(requests)-1, lockReq) + case tikvrpc.CmdPessimisticRollback: + rollbackReq := req.PessimisticRollback() + keys := make([][]byte, len(rollbackReq.Keys)) + for i, key := range rollbackReq.Keys { + keys[i] = append([]byte(nil), key...) + } + requests = append(requests, requestSummary{cmd: req.Type, keys: keys}) + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticRollbackResponse{}}, nil + default: + require.FailNow(t, "unexpected RPC", "command: %s", req.Type) + return nil, nil + } + } + return txn, &requests + } + + lockSharedKey := func(t *testing.T, txn *testTxn, primaryKey, upgradeKey []byte) { + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + require.NoError(t, txn.lockKeys(context.TODO(), lockCtx, nil, primaryKey)) + lockCtx.InShareMode = true + require.NoError(t, txn.lockKeys(context.TODO(), lockCtx, nil, upgradeKey)) + flags, err := txn.GetMemBuffer().GetFlags(upgradeKey) + require.NoError(t, err) + require.True(t, flags.HasLocked()) + require.True(t, flags.HasLockedInShareMode()) + } + + keysAsStrings := func(keys [][]byte) []string { + ret := make([]string, len(keys)) + for i, key := range keys { + ret[i] = string(key) + } + return ret + } + + t.Run("GateOffRejectsLocally", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + txn, requests := newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + + requestCount := len(*requests) + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + err := txn.lockKeys(context.TODO(), lockCtx, nil, upgradeKey) + require.ErrorContains(t, err, "upgrading a shared lock to an exclusive lock is not supported") + require.Len(t, *requests, requestCount) + + flags, err := txn.GetMemBuffer().GetFlags(upgradeKey) + require.NoError(t, err) + require.True(t, flags.HasLockedInShareMode()) + }) + + t.Run("GateOnSendsUpgradeSeparatelyAndPromotesLocalFlags", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + normalKey1 := []byte("normal-key-1") + normalKey2 := []byte("normal-key-2") + txn, requests := newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + *requests = (*requests)[:0] + + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + lockCtx.AllowSharedLockUpgrade = true + require.NoError(t, txn.lockKeys(context.TODO(), lockCtx, nil, normalKey1, upgradeKey, normalKey2)) + + require.Len(t, *requests, 2) + require.Equal(t, kvrpcpb.Op_PessimisticLock, (*requests)[0].op) + require.ElementsMatch(t, []string{string(normalKey1), string(normalKey2)}, keysAsStrings((*requests)[0].keys)) + require.Equal(t, kvrpcpb.Op_PessimisticLock, (*requests)[1].op) + require.Equal(t, []string{string(upgradeKey)}, keysAsStrings((*requests)[1].keys)) + + flags, err := txn.GetMemBuffer().GetFlags(upgradeKey) + require.NoError(t, err) + require.True(t, flags.HasLocked()) + require.False(t, flags.HasLockedInShareMode()) + + committer, err := TxnProbe{KVTxn: txn.KVTxn}.NewCommitter(1) + require.NoError(t, err) + mutations := committer.MutationsOfKeys([][]byte{upgradeKey}) + prewriteReq := committer.BuildPrewriteRequest(1, 1, 1, mutations, 1).Req.(*kvrpcpb.PrewriteRequest) + require.Len(t, prewriteReq.Mutations, 1) + require.Equal(t, kvrpcpb.Op_Lock, prewriteReq.Mutations[0].Op) + }) + + t.Run("RejectUpgradeInAggressiveLockingMode", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + txn, requests := newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + + requestCount := len(*requests) + txn.StartAggressiveLocking() + defer txn.CancelAggressiveLocking(context.Background()) + + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + lockCtx.AllowSharedLockUpgrade = true + err := txn.lockKeys(context.TODO(), lockCtx, nil, upgradeKey) + require.ErrorContains(t, err, "shared lock upgrade is not supported in aggressive/fair locking mode") + require.Len(t, *requests, requestCount) + require.True(t, txn.IsInAggressiveLockingMode()) + }) + + t.Run("UpgradeLockOnlyIfExistsRequiresReturnValues", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + txn, _ := newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + lockCtx.AllowSharedLockUpgrade = true + lockCtx.LockOnlyIfExists = true + err := txn.lockKeys(context.TODO(), lockCtx, nil, upgradeKey) + var noReturnValueErr *tikverr.ErrLockOnlyIfExistsNoReturnValue + require.ErrorAs(t, err, &noReturnValueErr) + require.Equal(t, upgradeKey, noReturnValueErr.LockKey) + }) + + t.Run("ExplicitUpgradeFailureKeepsExistingSharedHolderAndEarlierExclusiveLocks", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + normalKey := []byte("normal-key") + txn, requests := newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + if len(req.Mutations) == 1 && + req.Mutations[0].Op == kvrpcpb.Op_PessimisticLock && + string(req.Mutations[0].Key) == string(upgradeKey) { + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{ + Errors: []*kvrpcpb.KeyError{{ + Conflict: &kvrpcpb.WriteConflict{ + StartTs: 1, + ConflictTs: 2, + ConflictCommitTs: 3, + Key: upgradeKey, + Reason: kvrpcpb.WriteConflict_PessimisticRetry, + }, + }}, + }}, nil + } + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + *requests = (*requests)[:0] + + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + lockCtx.AllowSharedLockUpgrade = true + err := txn.lockKeys(context.TODO(), lockCtx, nil, normalKey, upgradeKey) + require.True(t, tikverr.IsErrWriteConflict(err)) + require.Len(t, *requests, 2) + require.Equal(t, []string{string(normalKey)}, keysAsStrings((*requests)[0].keys)) + require.Equal(t, []string{string(upgradeKey)}, keysAsStrings((*requests)[1].keys)) + + upgradeFlags, err := txn.GetMemBuffer().GetFlags(upgradeKey) + require.NoError(t, err) + require.True(t, upgradeFlags.HasLocked()) + require.True(t, upgradeFlags.HasLockedInShareMode()) + + normalFlags, err := txn.GetMemBuffer().GetFlags(normalKey) + require.NoError(t, err) + require.True(t, normalFlags.HasLocked()) + require.False(t, normalFlags.HasLockedInShareMode()) + + require.ElementsMatch(t, + []string{string(primaryKey), string(upgradeKey), string(normalKey)}, + keysAsStrings(TxnProbe{KVTxn: txn.KVTxn}.CollectLockedKeys())) + }) + + t.Run("LockUpgradeConflictReturnsTypedErrorWithoutRetrySemantics", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + txn, requests := newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + if len(req.Mutations) == 1 && + req.Mutations[0].Op == kvrpcpb.Op_PessimisticLock && + string(req.Mutations[0].Key) == string(upgradeKey) { + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{ + Errors: []*kvrpcpb.KeyError{{ + LockUpgradeConflict: &kvrpcpb.LockUpgradeConflict{ + Key: upgradeKey, + StartTs: 1, + OwnerStartTs: 2, + Reason: kvrpcpb.LockUpgradeConflict_SecondUpgrader, + }, + }}, + }}, nil + } + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + *requests = (*requests)[:0] + + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + lockCtx.AllowSharedLockUpgrade = true + err := txn.lockKeys(context.TODO(), lockCtx, nil, upgradeKey) + require.Error(t, err) + require.Len(t, *requests, 1) + require.Equal(t, []string{string(upgradeKey)}, keysAsStrings((*requests)[0].keys)) + require.False(t, tikverr.IsErrWriteConflict(err)) + require.False(t, tikverr.IsErrorUndetermined(err)) + + var retryable *tikverr.ErrRetryable + require.False(t, stderrs.As(err, &retryable)) + var deadlock *tikverr.ErrDeadlock + require.False(t, stderrs.As(err, &deadlock)) + + var conflict *tikverr.ErrLockUpgradeConflict + require.ErrorAs(t, err, &conflict) + require.Equal(t, []byte("upgrade-key"), conflict.Key) + require.Equal(t, uint64(1), conflict.StartTs) + require.Equal(t, uint64(2), conflict.OwnerStartTs) + require.Equal(t, kvrpcpb.LockUpgradeConflict_SecondUpgrader, conflict.Reason) + + flags, getErr := txn.GetMemBuffer().GetFlags(upgradeKey) + require.NoError(t, getErr) + require.True(t, flags.HasLocked()) + require.True(t, flags.HasLockedInShareMode()) + }) + + t.Run("SharedLockLostPoisonsTransactionButKeepsRollbackAvailable", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + var txn *testTxn + var requests *[]requestSummary + txn, requests = newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + if len(req.Mutations) == 1 && + req.Mutations[0].Op == kvrpcpb.Op_PessimisticLock && + string(req.Mutations[0].Key) == string(upgradeKey) { + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{ + Errors: []*kvrpcpb.KeyError{{ + SharedLockLost: &kvrpcpb.SharedLockLost{ + Key: append([]byte(nil), upgradeKey...), + StartTs: txn.StartTS(), + }, + }}, + }}, nil + } + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + *requests = (*requests)[:0] + + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + lockCtx.AllowSharedLockUpgrade = true + firstErr := txn.lockKeys(context.Background(), lockCtx, nil, upgradeKey) + require.Error(t, firstErr) + require.False(t, tikverr.IsErrorUndetermined(firstErr)) + require.Len(t, *requests, 1) + + var firstLost *tikverr.ErrSharedLockLost + require.ErrorAs(t, firstErr, &firstLost) + require.Equal(t, upgradeKey, firstLost.Key) + require.Equal(t, txn.StartTS(), firstLost.StartTs) + require.True(t, txn.Valid()) + + flags, err := txn.GetMemBuffer().GetFlags(upgradeKey) + require.NoError(t, err) + require.True(t, flags.HasLocked()) + require.True(t, flags.HasLockedInShareMode()) + + secondFatalErr := errors.WithStack(&tikverr.ErrSharedLockLost{ + SharedLockLost: &kvrpcpb.SharedLockLost{ + Key: []byte("second-upgrade-key"), + StartTs: txn.StartTS(), + }, + }) + txn.committer.setFatalTxnErr(secondFatalErr) + + requestCount := len(*requests) + laterErr := txn.LockKeys( + context.Background(), + kv.NewLockCtx(3, kv.LockNoWait, time.Now()), + []byte("later-key"), + ) + var laterLost *tikverr.ErrSharedLockLost + require.ErrorAs(t, laterErr, &laterLost) + require.Same(t, firstLost, laterLost) + require.Len(t, *requests, requestCount) + + commitErr := txn.Commit(context.Background()) + var commitLost *tikverr.ErrSharedLockLost + require.ErrorAs(t, commitErr, &commitLost) + require.Same(t, firstLost, commitLost) + require.Len(t, *requests, requestCount) + require.True(t, txn.Valid(), "fatal commit rejection must leave rollback available") + + require.NoError(t, txn.Rollback()) + require.False(t, txn.Valid()) + rolledBack := make([][]byte, 0, 2) + for _, request := range *requests { + if request.cmd == tikvrpc.CmdPessimisticRollback { + rolledBack = append(rolledBack, request.keys...) + } + } + require.ElementsMatch(t, [][]byte{primaryKey, upgradeKey}, rolledBack) + }) + + t.Run("LegacyAbortUpgradeFailureRemainsUndetermined", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + txn, requests := newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + if len(req.Mutations) == 1 && + req.Mutations[0].Op == kvrpcpb.Op_PessimisticLock && + string(req.Mutations[0].Key) == string(upgradeKey) { + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{ + Errors: []*kvrpcpb.KeyError{{ + Abort: "PessimisticLockNotFound during shared lock upgrade", + }}, + }}, nil + } + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + *requests = (*requests)[:0] + + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + lockCtx.AllowSharedLockUpgrade = true + firstErr := txn.lockKeys(context.Background(), lockCtx, nil, upgradeKey) + require.Error(t, firstErr) + require.True(t, tikverr.IsErrorUndetermined(firstErr)) + require.Len(t, *requests, 1) + + requestCount := len(*requests) + laterErr := txn.LockKeys( + context.Background(), + kv.NewLockCtx(3, kv.LockNoWait, time.Now()), + []byte("later-key"), + ) + require.Error(t, laterErr) + require.True(t, tikverr.IsErrorUndetermined(laterErr)) + require.Len(t, *requests, requestCount) + + commitErr := txn.Commit(context.Background()) + require.Error(t, commitErr) + require.True(t, tikverr.IsErrorUndetermined(commitErr)) + require.Len(t, *requests, requestCount) + }) + + t.Run("EmptyKeyErrorUpgradeFailureRemainsUndetermined", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + txn, requests := newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + if len(req.Mutations) == 1 && + req.Mutations[0].Op == kvrpcpb.Op_PessimisticLock && + string(req.Mutations[0].Key) == string(upgradeKey) { + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{ + Errors: []*kvrpcpb.KeyError{{}}, + }}, nil + } + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + *requests = (*requests)[:0] + + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + lockCtx.AllowSharedLockUpgrade = true + firstErr := txn.lockKeys(context.Background(), lockCtx, nil, upgradeKey) + require.Error(t, firstErr) + require.True(t, tikverr.IsErrorUndetermined(firstErr)) + require.Len(t, *requests, 1) + + requestCount := len(*requests) + laterErr := txn.LockKeys( + context.Background(), + kv.NewLockCtx(3, kv.LockNoWait, time.Now()), + []byte("later-key"), + ) + require.Error(t, laterErr) + require.True(t, tikverr.IsErrorUndetermined(laterErr)) + require.Len(t, *requests, requestCount) + + commitErr := txn.Commit(context.Background()) + require.Error(t, commitErr) + require.True(t, tikverr.IsErrorUndetermined(commitErr)) + require.Len(t, *requests, requestCount) + }) + + t.Run("UndeterminedTakesPrecedenceOverFatalState", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + txn, requests := newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + *requests = (*requests)[:0] + + fatalErr := errors.WithStack(&tikverr.ErrSharedLockLost{ + SharedLockLost: &kvrpcpb.SharedLockLost{ + Key: []byte("upgrade-key"), + StartTs: txn.StartTS(), + }, + }) + txn.committer.setFatalTxnErr(fatalErr) + txn.committer.setUndeterminedErr(errors.New("unknown upgrade outcome")) + + laterErr := txn.LockKeys( + context.Background(), + kv.NewLockCtx(3, kv.LockNoWait, time.Now()), + []byte("later-key"), + ) + require.Error(t, laterErr) + require.True(t, tikverr.IsErrorUndetermined(laterErr)) + require.Empty(t, *requests) + + commitErr := txn.Commit(context.Background()) + require.Error(t, commitErr) + require.True(t, tikverr.IsErrorUndetermined(commitErr)) + require.Empty(t, *requests) + require.False(t, txn.Valid()) + + require.ErrorIs(t, txn.Rollback(), tikverr.ErrInvalidTxn) + }) + + t.Run("OutcomeUnknownUpgradeFailureIsTransactionFatal", func(t *testing.T) { + primaryKey := []byte("primary-key") + upgradeKey := []byte("upgrade-key") + txn, requests := newRecorderTxn(t, func(callIndex int, req *kvrpcpb.PessimisticLockRequest) (*tikvrpc.Response, error) { + if len(req.Mutations) == 1 && + req.Mutations[0].Op == kvrpcpb.Op_PessimisticLock && + string(req.Mutations[0].Key) == string(upgradeKey) { + return &tikvrpc.Response{}, nil + } + return &tikvrpc.Response{Resp: &kvrpcpb.PessimisticLockResponse{}}, nil + }) + lockSharedKey(t, txn, primaryKey, upgradeKey) + *requests = (*requests)[:0] + + lockCtx := kv.NewLockCtx(2, kv.LockNoWait, time.Now()) + lockCtx.AllowSharedLockUpgrade = true + err := txn.lockKeys(context.TODO(), lockCtx, nil, upgradeKey) + require.Error(t, err) + require.True(t, tikverr.IsErrorUndetermined(err)) + require.Len(t, *requests, 1) + require.Equal(t, []string{string(upgradeKey)}, keysAsStrings((*requests)[0].keys)) + + flags, err := txn.GetMemBuffer().GetFlags(upgradeKey) + require.NoError(t, err) + require.True(t, flags.HasLocked()) + require.True(t, flags.HasLockedInShareMode()) + + err = txn.LockKeys(context.TODO(), kv.NewLockCtx(2, kv.LockNoWait, time.Now()), []byte("later-key")) + require.Error(t, err) + require.True(t, tikverr.IsErrorUndetermined(err)) + + err = txn.Commit(context.TODO()) + require.Error(t, err) + require.True(t, tikverr.IsErrorUndetermined(err)) + }) +} + func TestSharedLockCommitterIncompatibilities(t *testing.T) { t.Run("RejectSharedLockPrimaryKey", func(t *testing.T) { key := []byte("shared-key") diff --git a/util/redact/redact.go b/util/redact/redact.go index c56977631d..1ec1e7208d 100644 --- a/util/redact/redact.go +++ b/util/redact/redact.go @@ -92,6 +92,16 @@ func RedactKeyErrIfNecessary(err *kvrpcpb.KeyError) { e.Primary = redactMarker } } + if e := err.SharedLockLost; e != nil { + if len(e.Key) > 0 { + e.Key = redactMarker + } + } + if e := err.LockUpgradeConflict; e != nil { + if len(e.Key) > 0 { + e.Key = redactMarker + } + } if e := err.AlreadyExist; e != nil { if len(e.Key) > 0 { e.Key = redactMarker diff --git a/util/redact/redact_test.go b/util/redact/redact_test.go new file mode 100644 index 0000000000..96ef656379 --- /dev/null +++ b/util/redact/redact_test.go @@ -0,0 +1,35 @@ +package redact + +import ( + "testing" + + "github.com/pingcap/errors" + "github.com/pingcap/kvproto/pkg/kvrpcpb" + "github.com/stretchr/testify/require" +) + +func TestRedactKeyErrSharedLockLostAndLockUpgradeConflict(t *testing.T) { + originalMode := errors.RedactLogEnabled.Load() + t.Cleanup(func() { errors.RedactLogEnabled.Store(originalMode) }) + + newKeyErr := func() *kvrpcpb.KeyError { + return &kvrpcpb.KeyError{ + SharedLockLost: &kvrpcpb.SharedLockLost{Key: []byte("lost-key")}, + LockUpgradeConflict: &kvrpcpb.LockUpgradeConflict{ + Key: []byte("conflict-key"), + }, + } + } + + errors.RedactLogEnabled.Store(errors.RedactLogDisable) + unredacted := newKeyErr() + RedactKeyErrIfNecessary(unredacted) + require.Equal(t, []byte("lost-key"), unredacted.SharedLockLost.Key) + require.Equal(t, []byte("conflict-key"), unredacted.LockUpgradeConflict.Key) + + errors.RedactLogEnabled.Store(errors.RedactLogEnable) + redacted := newKeyErr() + RedactKeyErrIfNecessary(redacted) + require.Equal(t, []byte("?"), redacted.SharedLockLost.Key) + require.Equal(t, []byte("?"), redacted.LockUpgradeConflict.Key) +}