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
1 change: 1 addition & 0 deletions integration_tests/2pc_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -738,6 +738,7 @@ func (s *testCommitterSuite) TestPessimisticTTL() {
s.GreaterOrEqual(msBeforeLockExpired, int64(100))

lr := s.store.NewLockResolver()
defer lr.Close()
bo := tikv.NewBackofferWithVars(context.Background(), 5000, nil)
status, err := lr.GetTxnStatus(bo, txn.StartTS(), key2, 0, txn.StartTS(), true, false, nil)
s.Nil(err)
Expand Down
54 changes: 41 additions & 13 deletions integration_tests/lock_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -257,6 +257,7 @@ func (s *testLockSuite) TestCheckTxnStatusTTL() {

bo := tikv.NewBackofferWithVars(context.Background(), int(transaction.PrewriteMaxBackoff.Load()), nil)
lr := s.store.NewLockResolver()
defer lr.Close()
callerStartTS, err := s.store.GetOracle().GetTimestamp(bo.GetCtx(), &oracle.Option{TxnScope: oracle.GlobalTxnScope})
s.Nil(err)

Expand All @@ -270,7 +271,9 @@ func (s *testLockSuite) TestCheckTxnStatusTTL() {
// Rollback the txn.
lock := s.mustGetLock([]byte("key"))

err = s.store.NewLockResolver().ForceResolveLock(context.Background(), lock)
lr2 := s.store.NewLockResolver()
defer lr2.Close()
err = lr2.ForceResolveLock(context.Background(), lock)
s.Nil(err)

// Check its status is rollbacked.
Expand Down Expand Up @@ -303,7 +306,9 @@ func (s *testLockSuite) TestTxnHeartBeat() {
s.Equal(newTTL, uint64(6666))

lock := s.mustGetLock([]byte("key"))
err = s.store.NewLockResolver().ForceResolveLock(context.Background(), lock)
lr := s.store.NewLockResolver()
defer lr.Close()
err = lr.ForceResolveLock(context.Background(), lock)
s.Nil(err)

newTTL, err = s.store.SendTxnHeartbeat(context.Background(), []byte("key"), txn.StartTS(), 6666)
Expand All @@ -325,6 +330,7 @@ func (s *testLockSuite) TestCheckTxnStatus() {

bo := tikv.NewBackofferWithVars(context.Background(), int(transaction.PrewriteMaxBackoff.Load()), nil)
resolver := s.store.NewLockResolver()
defer resolver.Close()
// Call getTxnStatus to check the lock status.
status, err := resolver.GetTxnStatus(bo, txn.StartTS(), []byte("key"), currentTS, currentTS, true, false, nil)
s.Nil(err)
Expand All @@ -348,15 +354,19 @@ func (s *testLockSuite) TestCheckTxnStatus() {
// Then call getTxnStatus again and check the lock status.
currentTS, err = o.GetTimestamp(context.Background(), &oracle.Option{TxnScope: oracle.GlobalTxnScope})
s.Nil(err)
status, err = s.store.NewLockResolver().GetTxnStatus(bo, txn.StartTS(), []byte("key"), currentTS, 0, true, false, nil)
lr := s.store.NewLockResolver()
defer lr.Close()
status, err = lr.GetTxnStatus(bo, txn.StartTS(), []byte("key"), currentTS, 0, true, false, nil)
s.Nil(err)
s.Equal(status.TTL(), uint64(0))
s.Equal(status.CommitTS(), uint64(0))
s.Equal(status.Action(), kvrpcpb.Action_NoAction)

// Call getTxnStatus on a committed transaction.
startTS, commitTS := s.putKV([]byte("a"), []byte("a"))
status, err = s.store.NewLockResolver().GetTxnStatus(bo, startTS, []byte("a"), currentTS, currentTS, true, false, nil)
lr2 := s.store.NewLockResolver()
defer lr2.Close()
status, err = lr2.GetTxnStatus(bo, startTS, []byte("a"), currentTS, currentTS, true, false, nil)
s.Nil(err)
s.Equal(status.TTL(), uint64(0))
s.Equal(status.CommitTS(), commitTS)
Expand All @@ -382,6 +392,7 @@ func (s *testLockSuite) TestCheckTxnStatusNoWait() {
s.Nil(err)
bo := tikv.NewBackofferWithVars(context.Background(), int(transaction.PrewriteMaxBackoff.Load()), nil)
resolver := s.store.NewLockResolver()
defer resolver.Close()

// Call getTxnStatus for the TxnNotFound case.
_, err = resolver.GetTxnStatus(bo, txn.StartTS(), []byte("key"), currentTS, currentTS, false, false, nil)
Expand Down Expand Up @@ -539,6 +550,7 @@ func (s *testLockSuite) TestBatchResolveLocks() {
s.Greater(msBeforeLockExpired, int64(0))

lr := s.store.NewLockResolver()
defer lr.Close()
bo := tikv.NewGcResolveLockMaxBackoffer(context.Background())
loc, err := s.store.GetRegionCache().LocateKey(bo, locks[0].Primary)
s.Nil(err)
Expand Down Expand Up @@ -585,19 +597,25 @@ func (s *testLockSuite) TestZeroMinCommitTS() {
s.Nil(failpoint.Disable("tikvclient/mockZeroCommitTS"))

lock := s.mustGetLock([]byte("key"))
expire, pushed, _, err := s.store.NewLockResolver().ResolveLocksForRead(bo, 0, []*txnkv.Lock{lock}, true)
lr := s.store.NewLockResolver()
defer lr.Close()
expire, pushed, _, err := lr.ResolveLocksForRead(bo, 0, []*txnkv.Lock{lock}, true)
s.Nil(err)
s.Len(pushed, 0)
s.Greater(expire, int64(0))

expire, pushed, _, err = s.store.NewLockResolver().ResolveLocksForRead(bo, math.MaxUint64, []*txnkv.Lock{lock}, true)
lr2 := s.store.NewLockResolver()
defer lr2.Close()
expire, pushed, _, err = lr2.ResolveLocksForRead(bo, math.MaxUint64, []*txnkv.Lock{lock}, true)
s.Nil(err)
s.Len(pushed, 1)
s.Equal(expire, int64(0))

// Clean up this test.
lock.TTL = uint64(0)
expire, err = s.store.NewLockResolver().ResolveLocks(bo, 0, []*txnkv.Lock{lock})
lr3 := s.store.NewLockResolver()
defer lr3.Close()
expire, err = lr3.ResolveLocks(bo, 0, []*txnkv.Lock{lock})
s.Nil(err)
s.Equal(expire, int64(0))
}
Expand Down Expand Up @@ -635,6 +653,7 @@ func (s *testLockSuite) TestCheckLocksFallenBackFromAsyncCommit() {
s.True(lock.UseAsyncCommit)
bo := tikv.NewBackoffer(context.Background(), getMaxBackoff)
lr := s.store.NewLockResolver()
defer lr.Close()
status, err := lr.GetTxnStatusFromLock(bo, lock, 0, false)
s.Nil(err)
s.Equal(txnlock.LockProbe{}.GetPrimaryKeyFromTxnStatus(status), []byte("fb1"))
Expand Down Expand Up @@ -667,7 +686,9 @@ func (s *testLockSuite) TestResolveTxnFallenBackFromAsyncCommit() {

resolveStarted := time.Now()
for {
expire, err := s.store.NewLockResolver().ResolveLocks(bo, 0, []*txnkv.Lock{lock})
lr := s.store.NewLockResolver()
expire, err := lr.ResolveLocks(bo, 0, []*txnkv.Lock{lock})
lr.Close()
s.Nil(err)
if expire == 0 {
break
Expand All @@ -693,7 +714,9 @@ func (s *testLockSuite) TestBatchResolveTxnFallenBackFromAsyncCommit() {
bo := tikv.NewBackoffer(context.Background(), getMaxBackoff)
loc, err := s.store.GetRegionCache().LocateKey(bo, []byte("fb1"))
s.Nil(err)
ok, err := s.store.NewLockResolver().BatchResolveLocks(bo, []*txnkv.Lock{lock}, loc.Region)
lr := s.store.NewLockResolver()
defer lr.Close()
ok, err := lr.BatchResolveLocks(bo, []*txnkv.Lock{lock}, loc.Region)
s.Nil(err)
s.True(ok)

Expand Down Expand Up @@ -854,6 +877,7 @@ func (s *testLockSuite) TestStartHeartBeatAfterLockingPrimary() {
// Check the TTL should have been updated
lr := s.store.NewLockResolver()
status, err := lr.LockResolver.GetTxnStatus(txn.StartTS(), 0, []byte("a"))
lr.Close()
s.Nil(err)
s.False(status.IsCommitted())
s.Greater(status.TTL(), uint64(600))
Expand All @@ -873,13 +897,15 @@ func (s *testLockSuite) TestStartHeartBeatAfterLockingPrimary() {
// The original primary key "a" should be rolled back because its TTL is not updated
lr = s.store.NewLockResolver()
status, err = lr.LockResolver.GetTxnStatus(txn.StartTS(), 0, []byte("a"))
lr.Close()
s.Nil(err)
s.False(status.IsCommitted())
s.Equal(status.TTL(), uint64(0))

// The TTL of the new primary lock should be updated.
lr = s.store.NewLockResolver()
status, err = lr.LockResolver.GetTxnStatus(txn.StartTS(), 0, []byte("c"))
lr.Close()
s.Nil(err)
s.False(status.IsCommitted())
s.Greater(status.TTL(), uint64(1200))
Expand Down Expand Up @@ -939,7 +965,9 @@ func (s *testLockSuite) TestResolveLocksForRead() {
// rolled back
startTS, _ = s.lockKey([]byte("k2"), []byte("v2"), []byte("k22"), []byte("v22"), 3000, false, false)
lock = s.mustGetLock([]byte("k22"))
err := s.store.NewLockResolver().ForceResolveLock(ctx, lock)
lr := s.store.NewLockResolver()
defer lr.Close()
err := lr.ForceResolveLock(ctx, lock)
s.Nil(err)
resolvedLocks = append(resolvedLocks, startTS)
lock = s.mustGetLock([]byte("k2"))
Expand Down Expand Up @@ -989,12 +1017,12 @@ func (s *testLockSuite) TestResolveLocksForRead() {
}

bo := tikv.NewBackoffer(context.Background(), getMaxBackoff)
lr := s.store.NewLockResolver()
defer lr.Close()
lr2 := s.store.NewLockResolver()
defer lr2.Close()

// Sleep for a while to make sure the async commit lock "k5" expires, so it could be resolve commit.
time.Sleep(500 * time.Millisecond)
msBeforeExpired, resolved, committed, err := lr.ResolveLocksForRead(bo, readStartTS, locks, false)
msBeforeExpired, resolved, committed, err := lr2.ResolveLocksForRead(bo, readStartTS, locks, false)
s.Nil(err)
s.Greater(msBeforeExpired, int64(0))
s.Equal(resolvedLocks, resolved)
Expand Down
1 change: 1 addition & 0 deletions integration_tests/shared_lock_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -334,6 +334,7 @@ func (s *testSharedLockSuite) TestGCSharedLock() {
s.True(txn3.GetCommitter().IsTTLRunning(), "txn3's TTL manager should still be running after sleep")

lr := s.store.NewLockResolver()
defer lr.Close()
bo := tikv.NewGcResolveLockMaxBackoffer(context.Background())
ttl, err := lr.ResolveLocks(bo, 0, locks)
s.Nil(err)
Expand Down
Loading
Loading