mirror of https://github.com/tikv/client-go.git
snap_test: replace pingcap/check with testify (#154)
Signed-off-by: disksing <i@disksing.com>
This commit is contained in:
parent
3976d914fe
commit
2f0893faf6
|
@ -35,207 +35,206 @@ package tikv_test
|
|||
import (
|
||||
"context"
|
||||
"math"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
. "github.com/pingcap/check"
|
||||
"github.com/pingcap/failpoint"
|
||||
"github.com/pingcap/tidb/store/mockstore/unistore"
|
||||
"github.com/stretchr/testify/suite"
|
||||
tikverr "github.com/tikv/client-go/v2/error"
|
||||
"github.com/tikv/client-go/v2/mockstore"
|
||||
"github.com/tikv/client-go/v2/tikv"
|
||||
)
|
||||
|
||||
func TestSnapshotFail(t *testing.T) {
|
||||
suite.Run(t, new(testSnapshotFailSuite))
|
||||
}
|
||||
|
||||
type testSnapshotFailSuite struct {
|
||||
suite.Suite
|
||||
store tikv.StoreProbe
|
||||
}
|
||||
|
||||
var _ = SerialSuites(&testSnapshotFailSuite{})
|
||||
|
||||
func (s *testSnapshotFailSuite) SetUpSuite(c *C) {
|
||||
func (s *testSnapshotFailSuite) SetupSuite() {
|
||||
client, pdClient, cluster, err := unistore.New("")
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
unistore.BootstrapWithSingleStore(cluster)
|
||||
store, err := tikv.NewTestTiKVStore(fpClient{Client: client}, pdClient, nil, nil, 0)
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
s.store = tikv.StoreProbe{KVStore: store}
|
||||
}
|
||||
|
||||
func (s *testSnapshotFailSuite) cleanup(c *C) {
|
||||
func (s *testSnapshotFailSuite) TearDownTest() {
|
||||
txn, err := s.store.Begin()
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
iter, err := txn.Iter([]byte(""), []byte(""))
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
for iter.Valid() {
|
||||
err = txn.Delete(iter.Key())
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
err = iter.Next()
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
}
|
||||
c.Assert(txn.Commit(context.TODO()), IsNil)
|
||||
s.Require().Nil(txn.Commit(context.TODO()))
|
||||
}
|
||||
|
||||
func (s *testSnapshotFailSuite) TestBatchGetResponseKeyError(c *C) {
|
||||
func (s *testSnapshotFailSuite) TestBatchGetResponseKeyError() {
|
||||
// Meaningless to test with tikv because it has a mock key error
|
||||
if *mockstore.WithTiKV {
|
||||
return
|
||||
}
|
||||
defer s.cleanup(c)
|
||||
|
||||
// Put two KV pairs
|
||||
txn, err := s.store.Begin()
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
err = txn.Set([]byte("k1"), []byte("v1"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = txn.Set([]byte("k2"), []byte("v2"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = txn.Commit(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
|
||||
c.Assert(failpoint.Enable("tikvclient/rpcBatchGetResult", `1*return("keyError")`), IsNil)
|
||||
s.Require().Nil(failpoint.Enable("tikvclient/rpcBatchGetResult", `1*return("keyError")`))
|
||||
defer func() {
|
||||
c.Assert(failpoint.Disable("tikvclient/rpcBatchGetResult"), IsNil)
|
||||
s.Require().Nil(failpoint.Disable("tikvclient/rpcBatchGetResult"))
|
||||
}()
|
||||
|
||||
txn, err = s.store.Begin()
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
res, err := toTiDBTxn(&txn).BatchGet(context.Background(), toTiDBKeys([][]byte{[]byte("k1"), []byte("k2")}))
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(res, DeepEquals, map[string][]byte{"k1": []byte("v1"), "k2": []byte("v2")})
|
||||
s.Nil(err)
|
||||
s.Equal(res, map[string][]byte{"k1": []byte("v1"), "k2": []byte("v2")})
|
||||
}
|
||||
|
||||
func (s *testSnapshotFailSuite) TestScanResponseKeyError(c *C) {
|
||||
func (s *testSnapshotFailSuite) TestScanResponseKeyError() {
|
||||
// Meaningless to test with tikv because it has a mock key error
|
||||
if *mockstore.WithTiKV {
|
||||
return
|
||||
}
|
||||
defer s.cleanup(c)
|
||||
|
||||
// Put two KV pairs
|
||||
txn, err := s.store.Begin()
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
err = txn.Set([]byte("k1"), []byte("v1"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = txn.Set([]byte("k2"), []byte("v2"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = txn.Set([]byte("k3"), []byte("v3"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = txn.Commit(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
|
||||
c.Assert(failpoint.Enable("tikvclient/rpcScanResult", `1*return("keyError")`), IsNil)
|
||||
s.Require().Nil(failpoint.Enable("tikvclient/rpcScanResult", `1*return("keyError")`))
|
||||
txn, err = s.store.Begin()
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
iter, err := txn.Iter([]byte("a"), []byte("z"))
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(iter.Key(), DeepEquals, []byte("k1"))
|
||||
c.Assert(iter.Value(), DeepEquals, []byte("v1"))
|
||||
c.Assert(iter.Next(), IsNil)
|
||||
c.Assert(iter.Key(), DeepEquals, []byte("k2"))
|
||||
c.Assert(iter.Value(), DeepEquals, []byte("v2"))
|
||||
c.Assert(iter.Next(), IsNil)
|
||||
c.Assert(iter.Key(), DeepEquals, []byte("k3"))
|
||||
c.Assert(iter.Value(), DeepEquals, []byte("v3"))
|
||||
c.Assert(iter.Next(), IsNil)
|
||||
c.Assert(iter.Valid(), IsFalse)
|
||||
c.Assert(failpoint.Disable("tikvclient/rpcScanResult"), IsNil)
|
||||
s.Nil(err)
|
||||
s.Equal(iter.Key(), []byte("k1"))
|
||||
s.Equal(iter.Value(), []byte("v1"))
|
||||
s.Nil(iter.Next())
|
||||
s.Equal(iter.Key(), []byte("k2"))
|
||||
s.Equal(iter.Value(), []byte("v2"))
|
||||
s.Nil(iter.Next())
|
||||
s.Equal(iter.Key(), []byte("k3"))
|
||||
s.Equal(iter.Value(), []byte("v3"))
|
||||
s.Nil(iter.Next())
|
||||
s.False(iter.Valid())
|
||||
s.Require().Nil(failpoint.Disable("tikvclient/rpcScanResult"))
|
||||
|
||||
c.Assert(failpoint.Enable("tikvclient/rpcScanResult", `1*return("keyError")`), IsNil)
|
||||
s.Require().Nil(failpoint.Enable("tikvclient/rpcScanResult", `1*return("keyError")`))
|
||||
txn, err = s.store.Begin()
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
iter, err = txn.Iter([]byte("k2"), []byte("k4"))
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(iter.Key(), DeepEquals, []byte("k2"))
|
||||
c.Assert(iter.Value(), DeepEquals, []byte("v2"))
|
||||
c.Assert(iter.Next(), IsNil)
|
||||
c.Assert(iter.Key(), DeepEquals, []byte("k3"))
|
||||
c.Assert(iter.Value(), DeepEquals, []byte("v3"))
|
||||
c.Assert(iter.Next(), IsNil)
|
||||
c.Assert(iter.Valid(), IsFalse)
|
||||
c.Assert(failpoint.Disable("tikvclient/rpcScanResult"), IsNil)
|
||||
s.Nil(err)
|
||||
s.Equal(iter.Key(), []byte("k2"))
|
||||
s.Equal(iter.Value(), []byte("v2"))
|
||||
s.Nil(iter.Next())
|
||||
s.Equal(iter.Key(), []byte("k3"))
|
||||
s.Equal(iter.Value(), []byte("v3"))
|
||||
s.Nil(iter.Next())
|
||||
s.False(iter.Valid())
|
||||
s.Require().Nil(failpoint.Disable("tikvclient/rpcScanResult"))
|
||||
}
|
||||
|
||||
func (s *testSnapshotFailSuite) TestRetryMaxTsPointGetSkipLock(c *C) {
|
||||
defer s.cleanup(c)
|
||||
func (s *testSnapshotFailSuite) TestRetryMaxTsPointGetSkipLock() {
|
||||
|
||||
// Prewrite k1 and k2 with async commit but don't commit them
|
||||
txn, err := s.store.Begin()
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
err = txn.Set([]byte("k1"), []byte("v1"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = txn.Set([]byte("k2"), []byte("v2"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
txn.SetEnableAsyncCommit(true)
|
||||
|
||||
c.Assert(failpoint.Enable("tikvclient/asyncCommitDoNothing", "return"), IsNil)
|
||||
c.Assert(failpoint.Enable("tikvclient/twoPCShortLockTTL", "return"), IsNil)
|
||||
s.Require().Nil(failpoint.Enable("tikvclient/asyncCommitDoNothing", "return"))
|
||||
s.Require().Nil(failpoint.Enable("tikvclient/twoPCShortLockTTL", "return"))
|
||||
committer, err := txn.NewCommitter(1)
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = committer.Execute(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(failpoint.Disable("tikvclient/twoPCShortLockTTL"), IsNil)
|
||||
s.Nil(err)
|
||||
s.Require().Nil(failpoint.Disable("tikvclient/twoPCShortLockTTL"))
|
||||
|
||||
snapshot := s.store.GetSnapshot(math.MaxUint64)
|
||||
getCh := make(chan []byte)
|
||||
go func() {
|
||||
// Sleep a while to make the TTL of the first txn expire, then we make sure we resolve lock by this get
|
||||
time.Sleep(200 * time.Millisecond)
|
||||
c.Assert(failpoint.Enable("tikvclient/beforeSendPointGet", "1*off->pause"), IsNil)
|
||||
s.Require().Nil(failpoint.Enable("tikvclient/beforeSendPointGet", "1*off->pause"))
|
||||
res, err := snapshot.Get(context.Background(), []byte("k2"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
getCh <- res
|
||||
}()
|
||||
// The get should be blocked by the failpoint. But the lock should have been resolved.
|
||||
select {
|
||||
case res := <-getCh:
|
||||
c.Errorf("too early %s", string(res))
|
||||
s.Fail("too early %s", string(res))
|
||||
case <-time.After(1 * time.Second):
|
||||
}
|
||||
|
||||
// Prewrite k1 and k2 again without committing them
|
||||
txn, err = s.store.Begin()
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
txn.SetEnableAsyncCommit(true)
|
||||
err = txn.Set([]byte("k1"), []byte("v3"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = txn.Set([]byte("k2"), []byte("v4"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
committer, err = txn.NewCommitter(1)
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = committer.Execute(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
|
||||
c.Assert(failpoint.Disable("tikvclient/beforeSendPointGet"), IsNil)
|
||||
s.Require().Nil(failpoint.Disable("tikvclient/beforeSendPointGet"))
|
||||
|
||||
// After disabling the failpoint, the get request should bypass the new locks and read the old result
|
||||
select {
|
||||
case res := <-getCh:
|
||||
c.Assert(res, DeepEquals, []byte("v2"))
|
||||
s.Equal(res, []byte("v2"))
|
||||
case <-time.After(1 * time.Second):
|
||||
c.Errorf("get timeout")
|
||||
s.Fail("get timeout")
|
||||
}
|
||||
}
|
||||
|
||||
func (s *testSnapshotFailSuite) TestRetryPointGetResolveTS(c *C) {
|
||||
defer s.cleanup(c)
|
||||
|
||||
func (s *testSnapshotFailSuite) TestRetryPointGetResolveTS() {
|
||||
txn, err := s.store.Begin()
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(txn.Set([]byte("k1"), []byte("v1")), IsNil)
|
||||
s.Require().Nil(err)
|
||||
s.Nil(txn.Set([]byte("k1"), []byte("v1")))
|
||||
err = txn.Set([]byte("k2"), []byte("v2"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
txn.SetEnableAsyncCommit(false)
|
||||
txn.SetEnable1PC(false)
|
||||
txn.SetCausalConsistency(true)
|
||||
|
||||
// Prewrite the lock without committing it
|
||||
c.Assert(failpoint.Enable("tikvclient/beforeCommit", `pause`), IsNil)
|
||||
s.Require().Nil(failpoint.Enable("tikvclient/beforeCommit", `pause`))
|
||||
ch := make(chan struct{})
|
||||
committer, err := txn.NewCommitter(1)
|
||||
c.Assert(committer.GetPrimaryKey(), DeepEquals, []byte("k1"))
|
||||
s.Equal(committer.GetPrimaryKey(), []byte("k1"))
|
||||
go func() {
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = committer.Execute(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
ch <- struct{}{}
|
||||
}()
|
||||
|
||||
|
@ -244,15 +243,15 @@ func (s *testSnapshotFailSuite) TestRetryPointGetResolveTS(c *C) {
|
|||
// Should get nothing with max version, and **not pushing forward minCommitTS** of the primary lock
|
||||
snapshot := s.store.GetSnapshot(math.MaxUint64)
|
||||
_, err = snapshot.Get(context.Background(), []byte("k2"))
|
||||
c.Assert(tikverr.IsErrNotFound(err), IsTrue)
|
||||
s.True(tikverr.IsErrNotFound(err))
|
||||
|
||||
initialCommitTS := committer.GetCommitTS()
|
||||
c.Assert(failpoint.Disable("tikvclient/beforeCommit"), IsNil)
|
||||
s.Require().Nil(failpoint.Disable("tikvclient/beforeCommit"))
|
||||
|
||||
<-ch
|
||||
// check the minCommitTS is not pushed forward
|
||||
snapshot = s.store.GetSnapshot(initialCommitTS)
|
||||
v, err := snapshot.Get(context.Background(), []byte("k2"))
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(v, DeepEquals, []byte("v2"))
|
||||
s.Nil(err)
|
||||
s.Equal(v, []byte("v2"))
|
||||
}
|
||||
|
|
|
@ -37,149 +37,153 @@ import (
|
|||
"fmt"
|
||||
"math"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
. "github.com/pingcap/check"
|
||||
"github.com/pingcap/errors"
|
||||
"github.com/pingcap/failpoint"
|
||||
"github.com/pingcap/kvproto/pkg/kvrpcpb"
|
||||
tikverr "github.com/tikv/client-go/v2/error"
|
||||
"github.com/stretchr/testify/suite"
|
||||
"github.com/tikv/client-go/v2/error"
|
||||
"github.com/tikv/client-go/v2/logutil"
|
||||
"github.com/tikv/client-go/v2/tikv"
|
||||
"github.com/tikv/client-go/v2/tikvrpc"
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
func TestSnapshot(t *testing.T) {
|
||||
suite.Run(t, new(testSnapshotSuite))
|
||||
}
|
||||
|
||||
type testSnapshotSuite struct {
|
||||
suite.Suite
|
||||
store tikv.StoreProbe
|
||||
prefix string
|
||||
rowNums []int
|
||||
}
|
||||
|
||||
var _ = SerialSuites(&testSnapshotSuite{})
|
||||
|
||||
func (s *testSnapshotSuite) SetUpSuite(c *C) {
|
||||
s.store = tikv.StoreProbe{KVStore: NewTestStore(c)}
|
||||
func (s *testSnapshotSuite) SetupSuite() {
|
||||
s.store = tikv.StoreProbe{KVStore: NewTestStoreT(s.T())}
|
||||
s.prefix = fmt.Sprintf("snapshot_%d", time.Now().Unix())
|
||||
s.rowNums = append(s.rowNums, 1, 100, 191)
|
||||
}
|
||||
|
||||
func (s *testSnapshotSuite) TearDownSuite(c *C) {
|
||||
txn := s.beginTxn(c)
|
||||
func (s *testSnapshotSuite) TearDownSuite() {
|
||||
txn := s.beginTxn()
|
||||
scanner, err := txn.Iter(encodeKey(s.prefix, ""), nil)
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(scanner, NotNil)
|
||||
s.Nil(err)
|
||||
s.NotNil(scanner)
|
||||
for scanner.Valid() {
|
||||
k := scanner.Key()
|
||||
err = txn.Delete(k)
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
scanner.Next()
|
||||
}
|
||||
err = txn.Commit(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
err = s.store.Close()
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
}
|
||||
|
||||
func (s *testSnapshotSuite) beginTxn(c *C) tikv.TxnProbe {
|
||||
func (s *testSnapshotSuite) beginTxn() tikv.TxnProbe {
|
||||
txn, err := s.store.Begin()
|
||||
c.Assert(err, IsNil)
|
||||
s.Require().Nil(err)
|
||||
return txn
|
||||
}
|
||||
|
||||
func (s *testSnapshotSuite) checkAll(keys [][]byte, c *C) {
|
||||
txn := s.beginTxn(c)
|
||||
func (s *testSnapshotSuite) checkAll(keys [][]byte) {
|
||||
txn := s.beginTxn()
|
||||
snapshot := txn.GetSnapshot()
|
||||
m, err := snapshot.BatchGet(context.Background(), keys)
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
|
||||
scan, err := txn.Iter(encodeKey(s.prefix, ""), nil)
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
cnt := 0
|
||||
for scan.Valid() {
|
||||
cnt++
|
||||
k := scan.Key()
|
||||
v := scan.Value()
|
||||
v2, ok := m[string(k)]
|
||||
c.Assert(ok, IsTrue, Commentf("key: %q", k))
|
||||
c.Assert(v, BytesEquals, v2)
|
||||
s.True(ok, fmt.Sprintf("key: %q", k))
|
||||
s.Equal(v, v2)
|
||||
scan.Next()
|
||||
}
|
||||
err = txn.Commit(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(m, HasLen, cnt)
|
||||
s.Nil(err)
|
||||
s.Len(m, cnt)
|
||||
}
|
||||
|
||||
func (s *testSnapshotSuite) deleteKeys(keys [][]byte, c *C) {
|
||||
txn := s.beginTxn(c)
|
||||
func (s *testSnapshotSuite) deleteKeys(keys [][]byte) {
|
||||
txn := s.beginTxn()
|
||||
for _, k := range keys {
|
||||
err := txn.Delete(k)
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
}
|
||||
err := txn.Commit(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
}
|
||||
|
||||
func (s *testSnapshotSuite) TestBatchGet(c *C) {
|
||||
func (s *testSnapshotSuite) TestBatchGet() {
|
||||
for _, rowNum := range s.rowNums {
|
||||
logutil.BgLogger().Debug("test BatchGet",
|
||||
zap.Int("length", rowNum))
|
||||
txn := s.beginTxn(c)
|
||||
txn := s.beginTxn()
|
||||
for i := 0; i < rowNum; i++ {
|
||||
k := encodeKey(s.prefix, s08d("key", i))
|
||||
err := txn.Set(k, valueBytes(i))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
}
|
||||
err := txn.Commit(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
|
||||
keys := makeKeys(rowNum, s.prefix)
|
||||
s.checkAll(keys, c)
|
||||
s.deleteKeys(keys, c)
|
||||
s.checkAll(keys)
|
||||
s.deleteKeys(keys)
|
||||
}
|
||||
}
|
||||
|
||||
type contextKey string
|
||||
|
||||
func (s *testSnapshotSuite) TestSnapshotCache(c *C) {
|
||||
txn := s.beginTxn(c)
|
||||
c.Assert(txn.Set([]byte("x"), []byte("x")), IsNil)
|
||||
c.Assert(txn.Delete([]byte("y")), IsNil) // store data is affected by other tests.
|
||||
c.Assert(txn.Commit(context.Background()), IsNil)
|
||||
func (s *testSnapshotSuite) TestSnapshotCache() {
|
||||
txn := s.beginTxn()
|
||||
s.Nil(txn.Set([]byte("x"), []byte("x")))
|
||||
s.Nil(txn.Delete([]byte("y"))) // store data is affected by othe)
|
||||
s.Nil(txn.Commit(context.Background()))
|
||||
|
||||
txn = s.beginTxn(c)
|
||||
txn = s.beginTxn()
|
||||
snapshot := txn.GetSnapshot()
|
||||
_, err := snapshot.BatchGet(context.Background(), [][]byte{[]byte("x"), []byte("y")})
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
|
||||
c.Assert(failpoint.Enable("tikvclient/snapshot-get-cache-fail", `return(true)`), IsNil)
|
||||
s.Nil(failpoint.Enable("tikvclient/snapshot-get-cache-fail", `return(true)`))
|
||||
ctx := context.WithValue(context.Background(), contextKey("TestSnapshotCache"), true)
|
||||
_, err = snapshot.Get(ctx, []byte("x"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
|
||||
_, err = snapshot.Get(ctx, []byte("y"))
|
||||
c.Assert(tikverr.IsErrNotFound(err), IsTrue)
|
||||
s.True(error.IsErrNotFound(err))
|
||||
|
||||
c.Assert(failpoint.Disable("tikvclient/snapshot-get-cache-fail"), IsNil)
|
||||
s.Nil(failpoint.Disable("tikvclient/snapshot-get-cache-fail"))
|
||||
}
|
||||
|
||||
func (s *testSnapshotSuite) TestBatchGetNotExist(c *C) {
|
||||
func (s *testSnapshotSuite) TestBatchGetNotExist() {
|
||||
for _, rowNum := range s.rowNums {
|
||||
logutil.BgLogger().Debug("test BatchGetNotExist",
|
||||
zap.Int("length", rowNum))
|
||||
txn := s.beginTxn(c)
|
||||
txn := s.beginTxn()
|
||||
for i := 0; i < rowNum; i++ {
|
||||
k := encodeKey(s.prefix, s08d("key", i))
|
||||
err := txn.Set(k, valueBytes(i))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
}
|
||||
err := txn.Commit(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
|
||||
keys := makeKeys(rowNum, s.prefix)
|
||||
keys = append(keys, []byte("noSuchKey"))
|
||||
s.checkAll(keys, c)
|
||||
s.deleteKeys(keys, c)
|
||||
s.checkAll(keys)
|
||||
s.deleteKeys(keys)
|
||||
}
|
||||
}
|
||||
|
||||
|
@ -192,55 +196,55 @@ func makeKeys(rowNum int, prefix string) [][]byte {
|
|||
return keys
|
||||
}
|
||||
|
||||
func (s *testSnapshotSuite) TestSkipLargeTxnLock(c *C) {
|
||||
func (s *testSnapshotSuite) TestSkipLargeTxnLock() {
|
||||
x := []byte("x_key_TestSkipLargeTxnLock")
|
||||
y := []byte("y_key_TestSkipLargeTxnLock")
|
||||
txn := s.beginTxn(c)
|
||||
c.Assert(txn.Set(x, []byte("x")), IsNil)
|
||||
c.Assert(txn.Set(y, []byte("y")), IsNil)
|
||||
txn := s.beginTxn()
|
||||
s.Nil(txn.Set(x, []byte("x")))
|
||||
s.Nil(txn.Set(y, []byte("y")))
|
||||
ctx := context.Background()
|
||||
committer, err := txn.NewCommitter(0)
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
committer.SetLockTTL(3000)
|
||||
c.Assert(committer.PrewriteAllMutations(ctx), IsNil)
|
||||
s.Nil(committer.PrewriteAllMutations(ctx))
|
||||
|
||||
txn1 := s.beginTxn(c)
|
||||
txn1 := s.beginTxn()
|
||||
// txn1 is not blocked by txn in the large txn protocol.
|
||||
_, err = txn1.Get(ctx, x)
|
||||
c.Assert(tikverr.IsErrNotFound(errors.Trace(err)), IsTrue)
|
||||
s.True(error.IsErrNotFound(errors.Trace(err)))
|
||||
|
||||
res, err := toTiDBTxn(&txn1).BatchGet(ctx, toTiDBKeys([][]byte{x, y, []byte("z")}))
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(res, HasLen, 0)
|
||||
s.Nil(err)
|
||||
s.Len(res, 0)
|
||||
|
||||
// Commit txn, check the final commit ts is pushed.
|
||||
committer.SetCommitTS(txn.StartTS() + 1)
|
||||
c.Assert(committer.CommitMutations(ctx), IsNil)
|
||||
s.Nil(committer.CommitMutations(ctx))
|
||||
status, err := s.store.GetLockResolver().GetTxnStatus(txn.StartTS(), 0, x)
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(status.IsCommitted(), IsTrue)
|
||||
c.Assert(status.CommitTS(), Greater, txn1.StartTS())
|
||||
s.Nil(err)
|
||||
s.True(status.IsCommitted())
|
||||
s.Greater(status.CommitTS(), txn1.StartTS())
|
||||
}
|
||||
|
||||
func (s *testSnapshotSuite) TestPointGetSkipTxnLock(c *C) {
|
||||
func (s *testSnapshotSuite) TestPointGetSkipTxnLock() {
|
||||
x := []byte("x_key_TestPointGetSkipTxnLock")
|
||||
y := []byte("y_key_TestPointGetSkipTxnLock")
|
||||
txn := s.beginTxn(c)
|
||||
c.Assert(txn.Set(x, []byte("x")), IsNil)
|
||||
c.Assert(txn.Set(y, []byte("y")), IsNil)
|
||||
txn := s.beginTxn()
|
||||
s.Nil(txn.Set(x, []byte("x")))
|
||||
s.Nil(txn.Set(y, []byte("y")))
|
||||
ctx := context.Background()
|
||||
committer, err := txn.NewCommitter(0)
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
committer.SetLockTTL(3000)
|
||||
c.Assert(committer.PrewriteAllMutations(ctx), IsNil)
|
||||
s.Nil(committer.PrewriteAllMutations(ctx))
|
||||
|
||||
snapshot := s.store.GetSnapshot(math.MaxUint64)
|
||||
start := time.Now()
|
||||
c.Assert(committer.GetPrimaryKey(), BytesEquals, x)
|
||||
s.Equal(committer.GetPrimaryKey(), x)
|
||||
// Point get secondary key. Shouldn't be blocked by the lock and read old data.
|
||||
_, err = snapshot.Get(ctx, y)
|
||||
c.Assert(tikverr.IsErrNotFound(errors.Trace(err)), IsTrue)
|
||||
c.Assert(time.Since(start), Less, 500*time.Millisecond)
|
||||
s.True(error.IsErrNotFound(errors.Trace(err)))
|
||||
s.Less(time.Since(start), 500*time.Millisecond)
|
||||
|
||||
// Commit the primary key
|
||||
committer.SetCommitTS(txn.StartTS() + 1)
|
||||
|
@ -250,18 +254,18 @@ func (s *testSnapshotSuite) TestPointGetSkipTxnLock(c *C) {
|
|||
start = time.Now()
|
||||
// Point get secondary key. Should read committed data.
|
||||
value, err := snapshot.Get(ctx, y)
|
||||
c.Assert(err, IsNil)
|
||||
c.Assert(value, BytesEquals, []byte("y"))
|
||||
c.Assert(time.Since(start), Less, 500*time.Millisecond)
|
||||
s.Nil(err)
|
||||
s.Equal(value, []byte("y"))
|
||||
s.Less(time.Since(start), 500*time.Millisecond)
|
||||
}
|
||||
|
||||
func (s *testSnapshotSuite) TestSnapshotThreadSafe(c *C) {
|
||||
txn := s.beginTxn(c)
|
||||
func (s *testSnapshotSuite) TestSnapshotThreadSafe() {
|
||||
txn := s.beginTxn()
|
||||
key := []byte("key_test_snapshot_threadsafe")
|
||||
c.Assert(txn.Set(key, []byte("x")), IsNil)
|
||||
s.Nil(txn.Set(key, []byte("x")))
|
||||
ctx := context.Background()
|
||||
err := txn.Commit(context.Background())
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
|
||||
snapshot := s.store.GetSnapshot(math.MaxUint64)
|
||||
var wg sync.WaitGroup
|
||||
|
@ -270,9 +274,9 @@ func (s *testSnapshotSuite) TestSnapshotThreadSafe(c *C) {
|
|||
go func() {
|
||||
for i := 0; i < 30; i++ {
|
||||
_, err := snapshot.Get(ctx, key)
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
_, err = snapshot.BatchGet(ctx, [][]byte{key, []byte("key_not_exist")})
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
}
|
||||
wg.Done()
|
||||
}()
|
||||
|
@ -280,7 +284,7 @@ func (s *testSnapshotSuite) TestSnapshotThreadSafe(c *C) {
|
|||
wg.Wait()
|
||||
}
|
||||
|
||||
func (s *testSnapshotSuite) TestSnapshotRuntimeStats(c *C) {
|
||||
func (s *testSnapshotSuite) TestSnapshotRuntimeStats() {
|
||||
reqStats := tikv.NewRegionRequestRuntimeStats()
|
||||
tikv.RecordRegionRequestRuntimeStats(reqStats.Stats, tikvrpc.CmdGet, time.Second)
|
||||
tikv.RecordRegionRequestRuntimeStats(reqStats.Stats, tikvrpc.CmdGet, time.Millisecond)
|
||||
|
@ -290,11 +294,11 @@ func (s *testSnapshotSuite) TestSnapshotRuntimeStats(c *C) {
|
|||
snapshot.MergeRegionRequestStats(reqStats.Stats)
|
||||
bo := tikv.NewBackofferWithVars(context.Background(), 2000, nil)
|
||||
err := bo.BackoffWithMaxSleepTxnLockFast(30, errors.New("test"))
|
||||
c.Assert(err, IsNil)
|
||||
s.Nil(err)
|
||||
snapshot.RecordBackoffInfo(bo)
|
||||
snapshot.RecordBackoffInfo(bo)
|
||||
expect := "Get:{num_rpc:4, total_time:2s},txnLockFast_backoff:{num:2, total_time:60ms}"
|
||||
c.Assert(snapshot.FormatStats(), Equals, expect)
|
||||
s.Equal(snapshot.FormatStats(), expect)
|
||||
detail := &kvrpcpb.ExecDetailsV2{
|
||||
TimeDetail: &kvrpcpb.TimeDetail{
|
||||
WaitWallTimeMs: 100,
|
||||
|
@ -318,7 +322,7 @@ func (s *testSnapshotSuite) TestSnapshotRuntimeStats(c *C) {
|
|||
"rocksdb: {delete_skipped_count: 5, " +
|
||||
"key_skipped_count: 1, " +
|
||||
"block: {cache_hit_count: 10, read_count: 20, read_byte: 15 Bytes}}}"
|
||||
c.Assert(snapshot.FormatStats(), Equals, expect)
|
||||
s.Equal(snapshot.FormatStats(), expect)
|
||||
snapshot.MergeExecDetail(detail)
|
||||
expect = "Get:{num_rpc:4, total_time:2s},txnLockFast_backoff:{num:2, total_time:60ms}, " +
|
||||
"total_process_time: 200ms, total_wait_time: 200ms, " +
|
||||
|
@ -327,5 +331,5 @@ func (s *testSnapshotSuite) TestSnapshotRuntimeStats(c *C) {
|
|||
"rocksdb: {delete_skipped_count: 10, " +
|
||||
"key_skipped_count: 2, " +
|
||||
"block: {cache_hit_count: 20, read_count: 40, read_byte: 30 Bytes}}}"
|
||||
c.Assert(snapshot.FormatStats(), Equals, expect)
|
||||
s.Equal(snapshot.FormatStats(), expect)
|
||||
}
|
||||
|
|
Loading…
Reference in New Issue