-
Notifications
You must be signed in to change notification settings - Fork 268
Expand file tree
/
Copy pathresource_control_test.go
More file actions
145 lines (133 loc) · 5.33 KB
/
Copy pathresource_control_test.go
File metadata and controls
145 lines (133 loc) · 5.33 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
package resourcecontrol
import (
"testing"
"github.com/pingcap/kvproto/pkg/coprocessor"
"github.com/pingcap/kvproto/pkg/kvrpcpb"
"github.com/pingcap/kvproto/pkg/metapb"
"github.com/stretchr/testify/assert"
"github.com/tikv/client-go/v2/config"
"github.com/tikv/client-go/v2/tikvrpc"
)
func TestMakeRequestInfo(t *testing.T) {
// Test a non-write request.
req := &tikvrpc.Request{Req: &kvrpcpb.BatchGetRequest{}, Context: kvrpcpb.Context{Peer: &metapb.Peer{StoreId: 1}}}
info := MakeRequestInfo(req)
assert.False(t, info.IsWrite())
assert.Equal(t, uint64(0), info.WriteBytes())
assert.False(t, info.Bypass())
assert.Equal(t, uint64(1), info.StoreID())
// Test a prewrite request.
mutation := &kvrpcpb.Mutation{Key: []byte("foo"), Value: []byte("bar")}
prewriteReq := &kvrpcpb.PrewriteRequest{Mutations: []*kvrpcpb.Mutation{mutation}, PrimaryLock: []byte("baz")}
req = &tikvrpc.Request{Type: tikvrpc.CmdPrewrite, Req: prewriteReq, ReplicaNumber: 1, Context: kvrpcpb.Context{Peer: &metapb.Peer{StoreId: 2}}}
requestSource := "xxx_internal_others"
req.RequestSource = requestSource
info = MakeRequestInfo(req)
assert.True(t, info.IsWrite())
assert.Equal(t, uint64(9), info.WriteBytes())
assert.True(t, info.Bypass())
assert.Equal(t, uint64(2), info.StoreID())
// Test a commit request.
commitReq := &kvrpcpb.CommitRequest{Keys: [][]byte{[]byte("qux")}}
req = &tikvrpc.Request{Type: tikvrpc.CmdCommit, Req: commitReq, ReplicaNumber: 2, Context: kvrpcpb.Context{Peer: &metapb.Peer{StoreId: 3}}}
info = MakeRequestInfo(req)
assert.True(t, info.IsWrite())
assert.Equal(t, uint64(3), info.WriteBytes())
assert.False(t, info.Bypass())
assert.Equal(t, uint64(3), info.StoreID())
// Test Nil Peer in Context
req = &tikvrpc.Request{Type: tikvrpc.CmdCommit, Req: commitReq, ReplicaNumber: 2, Context: kvrpcpb.Context{}}
info = MakeRequestInfo(req)
assert.True(t, info.IsWrite())
assert.Equal(t, uint64(3), info.WriteBytes())
assert.False(t, info.Bypass())
assert.Equal(t, uint64(0), info.StoreID())
}
func TestMakeRequestInfoPredictedReadBytes(t *testing.T) {
// A read request may carry an optional PredictedReadBytes hint on
// tikvrpc.Request; MakeRequestInfo should propagate it to RequestInfo.
req := &tikvrpc.Request{
Req: &kvrpcpb.BatchGetRequest{},
Context: kvrpcpb.Context{Peer: &metapb.Peer{StoreId: 7}},
PredictedReadBytes: 256 * 1024,
}
info := MakeRequestInfo(req)
assert.False(t, info.IsWrite())
assert.Equal(t, uint64(256*1024), info.PredictedReadBytes(),
"predictedReadBytes should propagate from tikvrpc.Request")
// Without a hint, PredictedReadBytes defaults to 0. This test only
// checks client-go propagation; PD still decides request eligibility.
reqNoHint := &tikvrpc.Request{
Req: &kvrpcpb.BatchGetRequest{},
Context: kvrpcpb.Context{Peer: &metapb.Peer{StoreId: 7}},
}
infoNoHint := MakeRequestInfo(reqNoHint)
assert.Equal(t, uint64(0), infoNoHint.PredictedReadBytes(),
"zero hint on the request means zero on RequestInfo")
}
func TestMakeRequestInfoIsCop(t *testing.T) {
// Coprocessor requests must carry IsCop()==true so PD can scope
// paging_* metrics to them.
copReq := &tikvrpc.Request{
Type: tikvrpc.CmdCop,
Req: &coprocessor.Request{},
Context: kvrpcpb.Context{Peer: &metapb.Peer{StoreId: 1}},
}
assert.True(t, MakeRequestInfo(copReq).IsCop())
copStreamReq := &tikvrpc.Request{
Type: tikvrpc.CmdCopStream,
Req: &coprocessor.Request{},
Context: kvrpcpb.Context{Peer: &metapb.Peer{StoreId: 1}},
}
assert.True(t, MakeRequestInfo(copStreamReq).IsCop())
// Non-cop reads (Get, BatchGet, Scan) must carry IsCop()==false so
// PD ignores them in the paging accounting branch.
for _, req := range []*tikvrpc.Request{
{Type: tikvrpc.CmdGet, Req: &kvrpcpb.GetRequest{}, Context: kvrpcpb.Context{Peer: &metapb.Peer{StoreId: 1}}},
{Type: tikvrpc.CmdBatchGet, Req: &kvrpcpb.BatchGetRequest{}, Context: kvrpcpb.Context{Peer: &metapb.Peer{StoreId: 1}}},
{Type: tikvrpc.CmdScan, Req: &kvrpcpb.ScanRequest{}, Context: kvrpcpb.Context{Peer: &metapb.Peer{StoreId: 1}}},
} {
assert.False(t, MakeRequestInfo(req).IsCop(),
"non-cop cmd type %v must report IsCop()==false", req.Type)
}
}
func TestResponseInfoReadBytes(t *testing.T) {
resp := &tikvrpc.Response{
Resp: &coprocessor.Response{
ExecDetailsV2: &kvrpcpb.ExecDetailsV2{
ScanDetailV2: &kvrpcpb.ScanDetailV2{
TotalVersionsSize: 100,
ProcessedVersionsSize: 80,
RemoteTotalVersionsSize: 40,
RemoteProcessedVersionsSize: 30,
},
},
},
}
info := MakeResponseInfo(resp)
if config.NextGen {
assert.Equal(t, uint64(100), info.ReadBytes())
assert.Equal(t, uint64(40), info.RemoteReadBytes())
} else {
assert.Equal(t, uint64(80), info.ReadBytes())
assert.Zero(t, info.RemoteReadBytes())
}
if config.NextGen {
// Compatibility: when processed > total (older TiKV), use processed.
respCompat := &tikvrpc.Response{
Resp: &coprocessor.Response{
ExecDetailsV2: &kvrpcpb.ExecDetailsV2{
ScanDetailV2: &kvrpcpb.ScanDetailV2{
TotalVersionsSize: 80,
ProcessedVersionsSize: 100,
RemoteTotalVersionsSize: 60,
RemoteProcessedVersionsSize: 70,
},
},
},
}
infoCompat := MakeResponseInfo(respCompat)
assert.Equal(t, uint64(100), infoCompat.ReadBytes())
assert.Equal(t, uint64(70), infoCompat.RemoteReadBytes())
}
}