/
githubmirror
/
etcd
Обзор
Документация
Войти
/
githubmirror
/
etcd
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
tests/common/kv_test.go
372 строки
17 KB
Wei Fu
*: switch to new protobuf struct
15 май 2026, 17:59
15 май 2026, 17:59
df2b18b
Код
Авторство
О чём код?
// Copyright 2022 The etcd Authors // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // // http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // See the License for the specific language governing permissions and // limitations under the License. package common import ( "context" "fmt" "slices" "testing" "time" "github.com/google/go-cmp/cmp" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "google.golang.org/protobuf/proto" "google.golang.org/protobuf/testing/protocmp" "go.etcd.io/etcd/api/v3/etcdserverpb" "go.etcd.io/etcd/api/v3/mvccpb" clientv3 "go.etcd.io/etcd/client/v3" "go.etcd.io/etcd/server/v3/etcdserver/txn" "go.etcd.io/etcd/tests/v3/framework/config" "go.etcd.io/etcd/tests/v3/framework/interfaces" "go.etcd.io/etcd/tests/v3/framework/testutils" ) func TestKVPut(t *testing.T) { testRunner.BeforeTest(t) for _, tc := range clusterTestCases() { t.Run(tc.name, func(t *testing.T) { ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() clus := testRunner.NewCluster(ctx, t, config.WithClusterConfig(tc.config)) defer clus.Close() cc := testutils.MustClient(clus.Client()) testutils.ExecuteUntil(ctx, t, func() { key, value := "foo", "bar" _, err := cc.Put(ctx, key, value, config.PutOptions{}) require.NoErrorf(t, err, "count not put key %q", key) resp, err := cc.Get(ctx, key, config.GetOptions{}) require.NoErrorf(t, err, "count not get key %q, err: %s", key, err) assert.Lenf(t, resp.Kvs, 1, "Unexpected length of response, got %d", len(resp.Kvs)) assert.Equalf(t, string(resp.Kvs[0].Key), key, "Unexpected key, want %q, got %q", key, resp.Kvs[0].Key) assert.Equalf(t, string(resp.Kvs[0].Value), value, "Unexpected value, want %q, got %q", value, resp.Kvs[0].Value) }) }) } } func TestKVGet(t *testing.T) { testKVGet(t, false) } func TestKVGetStream(t *testing.T) { testKVGet(t, true) } func testKVGet(t *testing.T, stream bool) { testRunner.BeforeTest(t) for _, tc := range clusterTestCases() { t.Run(tc.name, func(t *testing.T) { ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) defer cancel() clus := testRunner.NewCluster(ctx, t, config.WithClusterConfig(tc.config)) defer clus.Close() cc := testutils.MustClient(clus.Client()) if stream && !clusterSupportsGetStream(ctx, t, clus) { t.Skip("RangeStream is not supported by this cluster") } testutils.ExecuteUntil(ctx, t, func() { resp, err := cc.Get(ctx, "", config.GetOptions{Prefix: true}) require.NoError(t, err) firstRev := resp.Header.Revision kvA := createKV("a", "aa1", firstRev+1, firstRev+1, 1) kvB := createKV("b", "a", firstRev+2, firstRev+2, 1) kvCV1 := createKV("c", "ac1", firstRev+3, firstRev+3, 1) kvCV2 := createKV("c", "ac2", firstRev+3, firstRev+4, 2) kvC := createKV("c", "aac", firstRev+3, firstRev+5, 3) kvFoo := createKV("foo", "bar", firstRev+6, firstRev+6, 1) kvFooAbc := createKV("foo/abc", "0", firstRev+7, firstRev+7, 1) kvFop := createKV("fop", "s", firstRev+8, firstRev+8, 1) inputs := []*mvccpb.KeyValue{kvA, kvB, kvCV1, kvCV2, kvC, kvFoo, kvFooAbc, kvFop} for i := range inputs { _, putError := cc.Put(ctx, string(inputs[i].Key), string(inputs[i].Value), config.PutOptions{}) require.NoErrorf(t, putError, "count not put key value %q", inputs[i]) } allKvs := []*mvccpb.KeyValue{kvA, kvB, kvC, kvFoo, kvFooAbc, kvFop} kvsByVersion := []*mvccpb.KeyValue{kvA, kvB, kvFoo, kvFooAbc, kvFop, kvC} reversedKvs := []*mvccpb.KeyValue{kvFop, kvFooAbc, kvFoo, kvC, kvB, kvA} kvsByValue := []*mvccpb.KeyValue{kvFooAbc, kvB, kvA, kvC, kvFoo, kvFop} kvsByValueDesc := []*mvccpb.KeyValue{kvFop, kvFoo, kvC, kvA, kvB, kvFooAbc} currentResp, err := cc.Get(ctx, "", config.GetOptions{Prefix: true}) require.NoError(t, err) currentHeader := &etcdserverpb.ResponseHeader{ ClusterId: currentResp.Header.ClusterId, Revision: currentResp.Header.Revision, } type testcase struct { name string begin string options config.GetOptions wantResponse *clientv3.GetResponse } tests := []testcase{ {name: "Get one specific key (a)", begin: "a", wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 1, Kvs: []*mvccpb.KeyValue{kvA}}}, {name: "Get one specific key (a), serializable", begin: "a", options: config.GetOptions{Serializable: true}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 1, Kvs: []*mvccpb.KeyValue{kvA}}}, {name: "Get [a, c)", begin: "a", options: config.GetOptions{End: "c"}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 2, Kvs: allKvs[:2]}}, {name: "blank key with --prefix option -> all KVs", begin: "", options: config.GetOptions{Prefix: true}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: allKvs}}, {name: "blank key with --from-key option -> all KVs", begin: "", options: config.GetOptions{FromKey: true}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: allKvs}}, {name: "Range covering all keys -> all KVs", begin: "a", options: config.GetOptions{End: "x"}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: allKvs}}, {name: "blank key with --prefix and revision -> [first key, entry at specified revision]", begin: "", options: config.GetOptions{Prefix: true, Revision: int(firstRev + 3)}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 3, Kvs: []*mvccpb.KeyValue{kvA, kvB, kvCV1}}}, {name: "--count-only for one single key", begin: "a", options: config.GetOptions{CountOnly: true}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 1, Kvs: nil}}, {name: "--prefix of foo -> all entries with the prefix", begin: "foo", options: config.GetOptions{Prefix: true}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 2, Kvs: allKvs[3:5]}}, {name: "--from-key of 'foo' -> <end>", begin: "foo", options: config.GetOptions{FromKey: true}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 3, Kvs: allKvs[3:]}}, {name: "blank key with limit set", begin: "", options: config.GetOptions{Prefix: true, Limit: 2}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: allKvs[:2], More: true}}, {name: "all kvs ordered by mod revision ascending", begin: "", options: config.GetOptions{Prefix: true, Order: clientv3.SortAscend, SortBy: clientv3.SortByModRevision}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: allKvs}}, {name: "all KVs ordered by version ascending", begin: "", options: config.GetOptions{Prefix: true, Order: clientv3.SortAscend, SortBy: clientv3.SortByVersion}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: kvsByVersion}}, {name: "all KVs ordered by key ascending, limit 2", begin: "", options: config.GetOptions{Prefix: true, Order: clientv3.SortAscend, SortBy: clientv3.SortByKey, Limit: 2}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: []*mvccpb.KeyValue{kvA, kvB}, More: true}}, {name: "range [b, z) ordered by key descending, limit 2", begin: "b", options: config.GetOptions{End: "z", Order: clientv3.SortDescend, SortBy: clientv3.SortByKey, Limit: 2}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 5, Kvs: []*mvccpb.KeyValue{kvFop, kvFooAbc}, More: true}}, {name: "all KVs ordered by create revision, unspecified sort order", begin: "", options: config.GetOptions{Prefix: true, Order: clientv3.SortNone, SortBy: clientv3.SortByCreateRevision}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: allKvs}}, {name: "all KVs ordered by create revision descending", begin: "", options: config.GetOptions{Prefix: true, Order: clientv3.SortDescend, SortBy: clientv3.SortByCreateRevision}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: reversedKvs}}, {name: "all KVs ordered by key descending", begin: "", options: config.GetOptions{Prefix: true, Order: clientv3.SortDescend, SortBy: clientv3.SortByKey}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: reversedKvs}}, {name: "all KVs ordered by value, unspecified sort order", begin: "", options: config.GetOptions{Prefix: true, Order: clientv3.SortNone, SortBy: clientv3.SortByValue}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: kvsByValue}}, {name: "all KVs ordered by value, ascending", begin: "", options: config.GetOptions{Prefix: true, Order: clientv3.SortAscend, SortBy: clientv3.SortByValue}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: kvsByValue}}, {name: "all KVs ordered by value descending", begin: "", options: config.GetOptions{Prefix: true, Order: clientv3.SortDescend, SortBy: clientv3.SortByValue}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: kvsByValueDesc}}, {name: "all KVs descending", begin: "", options: config.GetOptions{Prefix: true, Order: clientv3.SortDescend}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: reversedKvs}}, {name: "Get first version of 'c' by its revision", begin: "c", options: config.GetOptions{Revision: int(firstRev) + 3}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 1, Kvs: []*mvccpb.KeyValue{kvCV1}}}, {name: "Get second version of 'c' by its revision", begin: "c", options: config.GetOptions{Revision: int(firstRev) + 4}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 1, Kvs: []*mvccpb.KeyValue{kvCV2}}}, {name: "Get third version of 'c' by its revision", begin: "c", options: config.GetOptions{Revision: int(firstRev) + 5}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 1, Kvs: []*mvccpb.KeyValue{kvC}}}, {name: "Get the latest version of 'c'", begin: "c", wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 1, Kvs: []*mvccpb.KeyValue{kvC}}}, {name: "all KVs with mininum mod revision sorted by mod revision", begin: "", options: config.GetOptions{Prefix: true, MinModRevision: int(firstRev) + 3, SortBy: clientv3.SortByModRevision}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: allKvs[2:]}}, {name: "all KVs with maximum mod revision, sorted by key descending", begin: "", options: config.GetOptions{Prefix: true, MaxModRevision: int(firstRev) + 4, Order: clientv3.SortDescend, SortBy: clientv3.SortByKey}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: reversedKvs[4:]}}, {name: "all KVs with minimum create revision, sorted by version, descending", begin: "", options: config.GetOptions{Prefix: true, MinCreateRevision: int(firstRev) + 3, Order: clientv3.SortDescend, SortBy: clientv3.SortByVersion}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: allKvs[2:]}}, {name: "all KVs with maximimum create revision, sorted by value", begin: "", options: config.GetOptions{Prefix: true, MaxCreateRevision: int(firstRev) + 6, Order: clientv3.SortDescend, SortBy: clientv3.SortByValue}, wantResponse: &clientv3.GetResponse{Header: currentHeader, Count: 6, Kvs: kvsByValueDesc[1:5]}}, } testsWithKeysOnly := make([]testcase, 0, len(tests)) for _, otc := range tests { if otc.options.CountOnly { continue // can't use both --count-only and --keys-only at the same time } withKeysOnly := otc withKeysOnly.name = fmt.Sprintf("%s --keys-only", withKeysOnly.name) withKeysOnly.options.KeysOnly = true withKeysOnly.wantResponse = cloneGetResponseWithoutValues(otc.wantResponse) testsWithKeysOnly = append(testsWithKeysOnly, withKeysOnly) } for _, tt := range slices.Concat(tests, testsWithKeysOnly) { t.Run(tt.name, func(t *testing.T) { if stream && !rangeStreamSupports(tt.options) { t.Skip("options not supported by RangeStream") } opts := tt.options opts.Stream = stream resp, err := cc.Get(ctx, tt.begin, opts) require.NoErrorf(t, err, "count not get key %q, err: %s", tt.begin, err) resp.Header.MemberId = 0 resp.Header.RaftTerm = 0 assert.Emptyf(t, cmp.Diff( (*etcdserverpb.RangeResponse)(tt.wantResponse), (*etcdserverpb.RangeResponse)(resp), protocmp.Transform(), ), "-want, +got") }) } }) }) } } func createKV(key, val string, createRev, modRev, ver int64) *mvccpb.KeyValue { return &mvccpb.KeyValue{ Key: []byte(key), Value: []byte(val), CreateRevision: createRev, ModRevision: modRev, Version: ver, } } // clusterSupportsGetStream probes every cluster member with a RangeStream RPC and returns false if any member rejects it. func clusterSupportsGetStream(ctx context.Context, t *testing.T, clus interfaces.Cluster) bool { for _, m := range clus.Members() { _, err := m.Client().Get(ctx, "probe", config.GetOptions{Stream: true}) if err != nil { t.Logf("member does not support RangeStream: %v", err) return false } } return true } // rangeStreamSupports reports whether the server's RangeStream RPC accepts a // request with these options, mirroring v3rpc.checkRangeStreamRequest. func rangeStreamSupports(o config.GetOptions) bool { if !txn.IsDefaultOrdering( etcdserverpb.RangeRequest_SortTarget(o.SortBy), etcdserverpb.RangeRequest_SortOrder(o.Order), ) { return false } return !txn.HasRevisionFilters(&etcdserverpb.RangeRequest{ MinModRevision: int64(o.MinModRevision), MaxModRevision: int64(o.MaxModRevision), MinCreateRevision: int64(o.MinCreateRevision), MaxCreateRevision: int64(o.MaxCreateRevision), }) } func cloneGetResponseWithoutValues(resp *clientv3.GetResponse) *clientv3.GetResponse { clone := cloneGetResponse(resp) if clone == nil { return nil } for _, kv := range clone.Kvs { if kv != nil { kv.Value = nil } } return clone } func cloneGetResponse(resp *clientv3.GetResponse) *clientv3.GetResponse { if resp == nil { return nil } return (*clientv3.GetResponse)( proto.Clone( (*etcdserverpb.RangeResponse)(resp), ).(*etcdserverpb.RangeResponse), ) } func TestKVDelete(t *testing.T) { testRunner.BeforeTest(t) for _, tc := range clusterTestCases() { t.Run(tc.name, func(t *testing.T) { ctx, cancel := context.WithTimeout(t.Context(), 30*time.Second) defer cancel() clus := testRunner.NewCluster(ctx, t, config.WithClusterConfig(tc.config)) defer clus.Close() cc := testutils.MustClient(clus.Client()) testutils.ExecuteUntil(ctx, t, func() { kvs := []string{"a", "b", "c", "c/abc", "d"} tests := []struct { deleteKey string options config.DeleteOptions wantDeleted int wantKeys []string }{ { // delete all keys deleteKey: "", options: config.DeleteOptions{Prefix: true}, wantDeleted: 5, }, { // delete all keys deleteKey: "", options: config.DeleteOptions{FromKey: true}, wantDeleted: 5, }, { deleteKey: "a", options: config.DeleteOptions{End: "c"}, wantDeleted: 2, wantKeys: []string{"c", "c/abc", "d"}, }, { deleteKey: "c", wantDeleted: 1, wantKeys: []string{"a", "b", "c/abc", "d"}, }, { deleteKey: "c", options: config.DeleteOptions{Prefix: true}, wantDeleted: 2, wantKeys: []string{"a", "b", "d"}, }, { deleteKey: "c", options: config.DeleteOptions{FromKey: true}, wantDeleted: 3, wantKeys: []string{"a", "b"}, }, { deleteKey: "e", wantDeleted: 0, wantKeys: kvs, }, } for _, tt := range tests { for i := range kvs { _, err := cc.Put(ctx, kvs[i], "bar", config.PutOptions{}) require.NoErrorf(t, err, "count not put key %q", kvs[i]) } del, err := cc.Delete(ctx, tt.deleteKey, tt.options) require.NoErrorf(t, err, "count not get key %q, err", tt.deleteKey) assert.Equal(t, tt.wantDeleted, int(del.Deleted)) get, err := cc.Get(ctx, "", config.GetOptions{Prefix: true}) require.NoErrorf(t, err, "count not get key") kvs := testutils.KeysFromGetResponse(get) assert.Equal(t, tt.wantKeys, kvs) } }) }) } } func TestKVGetNoQuorum(t *testing.T) { testRunner.BeforeTest(t) tcs := []struct { name string options config.GetOptions wantError bool }{ { name: "Serializable", options: config.GetOptions{Serializable: true}, }, { name: "Linearizable", options: config.GetOptions{Serializable: false, Timeout: time.Second}, wantError: true, }, } for _, tc := range tcs { t.Run(tc.name, func(t *testing.T) { ctx, cancel := context.WithTimeout(t.Context(), 10*time.Second) defer cancel() clus := testRunner.NewCluster(ctx, t) defer clus.Close() clus.Members()[0].Stop() clus.Members()[1].Stop() cc := clus.Members()[2].Client() testutils.ExecuteUntil(ctx, t, func() { key := "foo" _, err := cc.Get(ctx, key, tc.options) if tc.wantError { require.Error(t, err) } else { require.NoError(t, err) } }) }) } }