From 01e0c44069be5f85ce5069f63f88b56ff51e950e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Wed, 24 Feb 2021 22:05:49 +0100 Subject: [PATCH 1/8] One shared etcd client --- etcd/embed.go | 276 +++++++++++--------------------------------------- 1 file changed, 58 insertions(+), 218 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index 159e92fcf..b80b7166f 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -36,6 +36,7 @@ import ( "go.etcd.io/etcd/clientv3/clientv3util" "go.etcd.io/etcd/clientv3/concurrency" "go.etcd.io/etcd/embed" + "go.etcd.io/etcd/etcdserver/api/v3client" "go.etcd.io/etcd/mvcc/mvccpb" "go.etcd.io/etcd/pkg/types" ) @@ -89,7 +90,8 @@ type Etcd struct { lm leaseMetadata - e *embed.Etcd + e *embed.Etcd + cli *clientv3.Client } func NewEtcd(opt Options, replicas int) *Etcd { @@ -103,6 +105,7 @@ func NewEtcd(opt Options, replicas int) *Etcd { // Close implements io.Closer func (e *Etcd) Close() error { _ = testhook.Closed(pilosa.NewAuditor(), e, nil) + if e.e != nil { if e.resizeCancel != nil { e.resizeCancel() @@ -114,6 +117,10 @@ func (e *Etcd) Close() error { <-e.e.Server.StopNotify() } + if e.cli != nil { + return e.cli.Close() + } + return nil } @@ -189,6 +196,7 @@ func (e *Etcd) Start(ctx context.Context) (disco.InitialClusterState, error) { } _ = testhook.Opened(pilosa.NewAuditor(), e, nil) e.e = etcd + e.cli = v3client.New(e.e.Server) select { case <-ctx.Done(): @@ -204,12 +212,6 @@ func (e *Etcd) Start(ctx context.Context) (disco.InitialClusterState, error) { } func (e *Etcd) startHeartbeat() error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "startHeartbeat: creates a new client") - } - defer cli.Close() - heartbeatID, ctx, heartbeatCancel, err := e.leaseKeepAlive(context.Background(), e.options.HeartbeatTTL) if err != nil { return errors.Wrap(err, "startHeartbeat: creates a new hearbeat") @@ -220,7 +222,7 @@ func (e *Etcd) startHeartbeat() error { value = disco.ClusterStateResizing } - if _, err := cli.Put(ctx, key, string(value), clientv3.WithLease(heartbeatID)); err != nil { + if _, err := e.cli.Put(ctx, key, string(value), clientv3.WithLease(heartbeatID)); err != nil { heartbeatCancel() return errors.Wrapf(err, "startHeartbeat: puts a key-value (%s, %s) with lease (%v)", key, value, heartbeatID) } @@ -231,17 +233,11 @@ func (e *Etcd) startHeartbeat() error { } func (e *Etcd) NodeState(ctx context.Context, peerID string) (disco.NodeState, error) { - cli, err := e.client() - if err != nil { - return disco.NodeStateUnknown, errors.Wrap(err, "NodeState: creates a new client") - } - defer cli.Close() - - return e.nodeState(ctx, cli, peerID) + return e.nodeState(ctx, peerID) } -func (e *Etcd) nodeState(ctx context.Context, cli *hookedClient, peerID string) (disco.NodeState, error) { - resp, err := cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly()) +func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, error) { + resp, err := e.cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly()) if err != nil { return disco.NodeStateUnknown, err } @@ -249,7 +245,7 @@ func (e *Etcd) nodeState(ctx context.Context, cli *hookedClient, peerID string) return disco.NodeStateResizing, nil } - resp, err = cli.Get(ctx, path.Join(heartbeatPrefix, peerID)) + resp, err = e.cli.Get(ctx, path.Join(heartbeatPrefix, peerID)) if err != nil { return disco.NodeStateUnknown, err } @@ -267,16 +263,9 @@ func (e *Etcd) nodeState(ctx context.Context, cli *hookedClient, peerID string) func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, error) { out := make(map[string]disco.NodeState) - - cli, err := e.client() - if err != nil { - return nil, errors.Wrap(err, "NodeStates") - } - defer cli.Close() - members := e.e.Server.Cluster().Members() for _, member := range members { - s, err := e.nodeState(ctx, cli, member.ID.String()) + s, err := e.nodeState(ctx, member.ID.String()) if err != nil { log.Println("NodeStates get node state", member.ID.String(), err.Error()) } @@ -287,15 +276,9 @@ func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, erro return out, nil } -func (e *Etcd) Started(ctx context.Context) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "Started") - } - defer cli.Close() - +func (e *Etcd) Started(ctx context.Context) (err error) { key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.NodeStateStarted - if _, err = cli.Put(ctx, key, string(value), clientv3.WithLease(e.heartbeatID)); err == nil { + if _, err = e.cli.Put(ctx, key, string(value), clientv3.WithLease(e.heartbeatID)); err == nil { e.lm.started = true } return err @@ -334,12 +317,6 @@ func (e *Etcd) ClusterState(ctx context.Context) (disco.ClusterState, error) { return disco.ClusterStateUnknown, nil } - cli, err := e.client() - if err != nil { - return disco.ClusterStateUnknown, errors.WithMessage(err, "ClusterState: creates a new client") - } - defer cli.Close() - var ( heartbeats int = 0 resize bool @@ -347,7 +324,7 @@ func (e *Etcd) ClusterState(ctx context.Context) (disco.ClusterState, error) { ) members := e.e.Server.Cluster().Members() for _, m := range members { - ns, err := e.nodeState(ctx, cli, m.ID.String()) + ns, err := e.nodeState(ctx, m.ID.String()) if err != nil { log.Println("ClusterState get node state", err.Error()) continue @@ -384,12 +361,6 @@ func (e *Etcd) ClusterState(ctx context.Context) (disco.ClusterState, error) { } func (e *Etcd) Resize(ctx context.Context) (func([]byte) error, error) { - cli, err := e.client() - if err != nil { - return nil, errors.Wrap(err, "Resize: creates a new client") - } - defer cli.Close() - resizeID, ctx, resizeCancel, err := e.leaseKeepAlive(ctx, e.options.HeartbeatTTL) if err != nil { return nil, errors.Wrap(err, "Resize: creates a new hearbeat") @@ -397,7 +368,7 @@ func (e *Etcd) Resize(ctx context.Context) (func([]byte) error, error) { // Check if key exists - maybe we are still resizing key := path.Join(resizePrefix, e.e.Server.ID().String()) - txnResp, err := cli.Txn(ctx). + txnResp, err := e.cli.Txn(ctx). If(clientv3util.KeyMissing(key)). Then(clientv3.OpPut(key, "", clientv3.WithLease(resizeID))). Commit() @@ -427,14 +398,8 @@ func (e *Etcd) DoneResize() error { } func (e *Etcd) Watch(ctx context.Context, peerID string, onUpdate func([]byte) error) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "Watch: creates a new client") - } - defer cli.Close() - key := path.Join(resizePrefix, peerID) - for resp := range cli.Watch(ctx, key) { + for resp := range e.cli.Watch(ctx, key) { if err := resp.Err(); err != nil { return errors.Wrapf(err, "Watch: key (%s) response", key) } @@ -464,13 +429,7 @@ func (e *Etcd) DeleteNode(ctx context.Context, nodeID string) error { return err } - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "DeleteNode: creates a new client") - } - defer cli.Close() - - _, err = cli.MemberRemove(ctx, uint64(id)) + _, err = e.cli.MemberRemove(ctx, uint64(id)) if err != nil { return errors.Wrap(err, "DeleteNode: removes an existing member from the cluster") } @@ -534,13 +493,7 @@ func (e *Etcd) Schema(ctx context.Context) (disco.Schema, error) { } func (e *Etcd) Metadata(ctx context.Context, peerID string) ([]byte, error) { - cli, err := e.client() - if err != nil { - return nil, errors.Wrap(err, "Metadata") - } - defer cli.Close() - - resp, err := cli.Get(ctx, path.Join(metadataPrefix, peerID)) + resp, err := e.cli.Get(ctx, path.Join(metadataPrefix, peerID)) if err != nil { return nil, err } @@ -569,12 +522,6 @@ func (e *Etcd) SetMetadata(ctx context.Context, metadata []byte) error { } func (e *Etcd) CreateIndex(ctx context.Context, name string, val []byte) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "CreateIndex: creating client") - } - defer cli.Close() - key := schemaPrefix + name // Set up Op to write index value as bytes. @@ -582,7 +529,7 @@ func (e *Etcd) CreateIndex(ctx context.Context, name string, val []byte) error { op.WithValueBytes(val) // Check for key existence, and execute Op within a transaction. - resp, err := cli.KV.Txn(ctx). + resp, err := e.cli.Txn(ctx). If(clientv3util.KeyMissing(key)). Then(op). Commit() @@ -601,16 +548,10 @@ func (e *Etcd) Index(ctx context.Context, name string) ([]byte, error) { return e.getKeyBytes(ctx, schemaPrefix+name) } -func (e *Etcd) DeleteIndex(ctx context.Context, name string) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "DeleteIndex: creating client") - } - defer cli.Close() - +func (e *Etcd) DeleteIndex(ctx context.Context, name string) (err error) { key := schemaPrefix + name // Deleting index and fields in one transaction. - _, err = cli.KV.Txn(ctx). + _, err = e.cli.Txn(ctx). If(clientv3.Compare(clientv3.Version(key), ">", -1)). Then( clientv3.OpDelete(key+"/", clientv3.WithPrefix()), // deleting index fields @@ -626,12 +567,6 @@ func (e *Etcd) Field(ctx context.Context, indexName string, name string) ([]byte } func (e *Etcd) CreateField(ctx context.Context, indexName string, name string, val []byte) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "CreateField: creating client") - } - defer cli.Close() - key := schemaPrefix + indexName + "/" + name // Set up Op to write field value as bytes. @@ -639,7 +574,7 @@ func (e *Etcd) CreateField(ctx context.Context, indexName string, name string, v op.WithValueBytes(val) // Check for key existence, and execute Op within a transaction. - resp, err := cli.KV.Txn(ctx). + resp, err := e.cli.Txn(ctx). If(clientv3util.KeyMissing(key)). Then(op). Commit() @@ -654,16 +589,10 @@ func (e *Etcd) CreateField(ctx context.Context, indexName string, name string, v return nil } -func (e *Etcd) DeleteField(ctx context.Context, indexname string, name string) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "DeleteField: creating client") - } - defer cli.Close() - +func (e *Etcd) DeleteField(ctx context.Context, indexname string, name string) (err error) { key := schemaPrefix + indexname + "/" + name // Deleting field and views in one transaction. - _, err = cli.KV.Txn(ctx). + _, err = e.cli.Txn(ctx). If(clientv3.Compare(clientv3.Version(key), ">", -1)). Then( clientv3.OpDelete(key+"/", clientv3.WithPrefix()), // deleting field views @@ -681,17 +610,11 @@ func (e *Etcd) View(ctx context.Context, indexName, fieldName, name string) (boo // CreateView differs from CreateIndex and CreateField in that it does not // return an error if the view already exists. If this logic needs to be // changed, we likely need to return disco.ErrViewExists. -func (e *Etcd) CreateView(ctx context.Context, indexName, fieldName, name string) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "CreateView: creating client") - } - defer cli.Close() - +func (e *Etcd) CreateView(ctx context.Context, indexName, fieldName, name string) (err error) { key := schemaPrefix + indexName + "/" + fieldName + "/" + name // Check for key existence, and execute Op within a transaction. - _, err = cli.KV.Txn(ctx). + _, err = e.cli.Txn(ctx). If(clientv3util.KeyMissing(key)). Then(clientv3.OpPut(key, "")). Commit() @@ -707,13 +630,7 @@ func (e *Etcd) DeleteView(ctx context.Context, indexName, fieldName, name string } func (e *Etcd) putKey(ctx context.Context, key, val string, opts ...clientv3.OpOption) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "putKey: creates a new client") - } - defer cli.Close() - - if _, err := cli.KV.Put(ctx, key, val, opts...); err != nil { + if _, err := e.cli.Put(ctx, key, val, opts...); err != nil { return errors.Wrapf(err, "putKey: Put(%s, %s)", key, val) } @@ -721,14 +638,8 @@ func (e *Etcd) putKey(ctx context.Context, key, val string, opts ...clientv3.OpO } func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { - cli, err := e.client() - if err != nil { - return nil, errors.Wrap(err, "getKeyBytes: creates a new client") - } - defer cli.Close() - // Get the current value for the key. - resp, err := cli.Get(ctx, key) + resp, err := e.cli.Get(ctx, key) if err != nil { return nil, err } @@ -742,13 +653,7 @@ func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { } func (e *Etcd) getKey(ctx context.Context, key string) ([]string, [][]byte, error) { - cli, err := e.client() - if err != nil { - return nil, nil, errors.Wrap(err, "getKey: creates a new client") - } - defer cli.Close() - - resp, err := cli.KV.Txn(ctx). + resp, err := e.cli.Txn(ctx). If(clientv3.Compare(clientv3.Version(key), ">", -1)). Then(clientv3.OpGet(key, clientv3.WithPrefix())). Commit() @@ -776,13 +681,7 @@ func (e *Etcd) getKey(ctx context.Context, key string) ([]string, [][]byte, erro } func (e *Etcd) keyExists(ctx context.Context, key string) (bool, error) { - cli, err := e.client() - if err != nil { - return false, errors.Wrap(err, "keyExists: creates a new client") - } - defer cli.Close() - - resp, err := cli.Get(ctx, key, clientv3.WithCountOnly()) + resp, err := e.cli.Get(ctx, key, clientv3.WithCountOnly()) if err != nil { return false, err } @@ -792,19 +691,13 @@ func (e *Etcd) keyExists(ctx context.Context, key string) (bool, error) { return false, nil } -func (e *Etcd) delKey(ctx context.Context, key string, withPrefix bool) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "delKey: creates a new client") - } - defer cli.Close() - +func (e *Etcd) delKey(ctx context.Context, key string, withPrefix bool) (err error) { var opts []clientv3.OpOption if withPrefix { opts = append(opts, clientv3.WithPrefix()) } - _, err = cli.KV.Txn(ctx). + _, err = e.cli.KV.Txn(ctx). If(clientv3.Compare(clientv3.Version(key), ">", -1)). Then(clientv3.OpDelete(key, opts...)). Commit() @@ -816,15 +709,8 @@ func (e *Etcd) delKey(ctx context.Context, key string, withPrefix bool) error { // then refreshes it periodically, and cancels it when done. it yields the lease ID, // and also a context and cancelfunc that can be used to abort the heartbeat. func (e *Etcd) leaseKeepAlive(ctx context.Context, ttl int64) (clientv3.LeaseID, context.Context, context.CancelFunc, error) { - cli, err := e.client() - if err != nil { - return 0, nil, nil, errors.Wrap(err, "leaseKeepAlive: creates a new client") - } - defer cli.Close() - ctx, cancelFunc := context.WithCancel(ctx) - - leaseResp, err := cli.Grant(ctx, ttl) + leaseResp, err := e.cli.Grant(ctx, ttl) if err != nil { cancelFunc() return 0, nil, nil, errors.Wrapf(err, "leaseKeepAlive: creates a new lease (TTL: %v)", ttl) @@ -842,28 +728,19 @@ func (e *Etcd) leaseKeepAlive(ctx context.Context, ttl int64) (clientv3.LeaseID, // here, resulting in massive piles of excess goroutines. revoker, cancel := context.WithTimeout(context.Background(), time.Duration(ttl)) defer cancel() - if cli, err := e.client(); err != nil { - log.Printf("leaseKeepAlive: creates a new client: %v\n", err) - } else { - if _, err := cli.Revoke(revoker, leaseResp.ID); err != nil { - log.Printf("leaseKeepAlive: revokes the lease (ID: %x): %#v\n", leaseResp.ID, err) - } - cli.Close() + + if _, err := e.cli.Revoke(revoker, leaseResp.ID); err != nil { + log.Printf("leaseKeepAlive: revokes the lease (ID: %x): %#v\n", leaseResp.ID, err) } return case <-ticker.C: - if cli, err := e.client(); err != nil { - log.Printf("leaseKeepAlive: creates a new client: %v\n", err) - } else { - if _, err = cli.KeepAliveOnce(ctx, leaseResp.ID); err != nil { - log.Printf("leaseKeepAlive: renews the lease (ID: %x): %v\n", leaseResp.ID, err) - } - cli.Close() + if _, err = e.cli.KeepAliveOnce(ctx, leaseResp.ID); err != nil { + log.Printf("leaseKeepAlive: renews the lease (ID: %x): %v\n", leaseResp.ID, err) } } } } - go keepaliveFunc(1 * time.Second) + go keepaliveFunc(time.Second) return leaseResp.ID, ctx, cancelFunc, nil } @@ -882,19 +759,6 @@ func (h *hookedClient) Close() { h.Client.Close() } -func (e *Etcd) client() (*hookedClient, error) { - urls := e.e.Server.Cluster().ClientURLs() - - cli, err := clientv3.NewFromURLs(urls) - if err != nil { - return nil, errors.Wrapf(err, "creates a new etcd client from URLs (%v)", urls) - } - - // Temporarily disabled, see comment in Close above. - // _ = testhook.Opened(pilosa.NewAuditor(), cli, nil) - return &hookedClient{Client: cli}, nil -} - func memberList(cli *hookedClient) (ids []uint64, names []string, urls []string) { ml, err := cli.MemberList(context.TODO()) if err != nil { @@ -922,20 +786,14 @@ func memberAdd(cli *hookedClient, peerURL string) (id uint64, name string) { // Shards implements the Sharder interface. func (e *Etcd) Shards(ctx context.Context, index, field string) (*roaring.Bitmap, error) { - cli, err := e.client() - if err != nil { - return nil, errors.Wrap(err, "Shards: creating client") - } - defer cli.Close() - - return e.shards(ctx, cli, index, field) + return e.shards(ctx, index, field) } -func (e *Etcd) shards(ctx context.Context, cli *hookedClient, index, field string) (*roaring.Bitmap, error) { +func (e *Etcd) shards(ctx context.Context, index, field string) (*roaring.Bitmap, error) { key := path.Join(shardPrefix, index, field) // Get the current shards for the field. - resp, err := cli.Get(ctx, key) + resp, err := e.cli.Get(ctx, key) if err != nil { return nil, err } @@ -956,12 +814,6 @@ func (e *Etcd) shards(ctx context.Context, cli *hookedClient, index, field strin // AddShards implements the Sharder interface. func (e *Etcd) AddShards(ctx context.Context, index, field string, shards *roaring.Bitmap) (*roaring.Bitmap, error) { - cli, err := e.client() - if err != nil { - return nil, errors.Wrap(err, "AddShards: creating client") - } - defer cli.Close() - key := path.Join(shardPrefix, index, field) // This tended to add more overhead than it saved. @@ -974,7 +826,7 @@ func (e *Etcd) AddShards(ctx context.Context, index, field string, shards *roari // } // Create a session to acquire a lock. - sess, _ := concurrency.NewSession(cli.Client) + sess, _ := concurrency.NewSession(e.cli) defer sess.Close() muKey := path.Join(lockPrefix, index, field) @@ -986,7 +838,7 @@ func (e *Etcd) AddShards(ctx context.Context, index, field string, shards *roari } // Read shards within lock. - globalShards, err := e.shards(ctx, cli, index, field) + globalShards, err := e.shards(ctx, index, field) if err != nil { return nil, errors.Wrap(err, "reading shards") } @@ -1003,7 +855,7 @@ func (e *Etcd) AddShards(ctx context.Context, index, field string, shards *roari op := clientv3.OpPut(key, "") op.WithValueBytes(buf.Bytes()) - if _, err := cli.Do(ctx, op); err != nil { + if _, err := e.cli.Do(ctx, op); err != nil { return nil, errors.Wrap(err, "doing op") } @@ -1017,17 +869,11 @@ func (e *Etcd) AddShards(ctx context.Context, index, field string, shards *roari // AddShard implements the Sharder interface. func (e *Etcd) AddShard(ctx context.Context, index, field string, shard uint64) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "AddShard: creating client") - } - defer cli.Close() - key := path.Join(shardPrefix, index, field) // Read shards outside of a lock just to check if shard is already included. // If shard is already included, no-op. - if shards, err := e.shards(ctx, cli, index, field); err != nil { + if shards, err := e.shards(ctx, index, field); err != nil { return errors.Wrap(err, "reading shards") } else if shards.Contains(shard) { return nil @@ -1039,7 +885,7 @@ func (e *Etcd) AddShard(ctx context.Context, index, field string, shard uint64) // write shards to etcd. // Create a session to acquire a lock. - sess, _ := concurrency.NewSession(cli.Client) + sess, _ := concurrency.NewSession(e.cli) defer sess.Close() muKey := path.Join(lockPrefix, index, field) @@ -1051,7 +897,7 @@ func (e *Etcd) AddShard(ctx context.Context, index, field string, shard uint64) } // Read shards again (within lock). - shards, err := e.shards(ctx, cli, index, field) + shards, err := e.shards(ctx, index, field) if err != nil { return errors.Wrap(err, "reading shards") } @@ -1072,7 +918,7 @@ func (e *Etcd) AddShard(ctx context.Context, index, field string, shard uint64) op := clientv3.OpPut(key, "") op.WithValueBytes(buf.Bytes()) - if _, err := cli.Do(ctx, op); err != nil { + if _, err := e.cli.Do(ctx, op); err != nil { return errors.Wrap(err, "doing op") } @@ -1086,17 +932,11 @@ func (e *Etcd) AddShard(ctx context.Context, index, field string, shard uint64) // RemoveShard implements the Sharder interface. func (e *Etcd) RemoveShard(ctx context.Context, index, field string, shard uint64) error { - cli, err := e.client() - if err != nil { - return errors.Wrap(err, "RemoveShard: creating client") - } - defer cli.Close() - key := path.Join(shardPrefix, index, field) // Read shards outside of a lock just to check if shard is already excluded. // If shard is already excluded, no-op. - if shards, err := e.shards(ctx, cli, index, field); err != nil { + if shards, err := e.shards(ctx, index, field); err != nil { return errors.Wrap(err, "reading shards") } else if !shards.Contains(shard) { return nil @@ -1108,7 +948,7 @@ func (e *Etcd) RemoveShard(ctx context.Context, index, field string, shard uint6 // write shards to etcd. // Create a session to acquire a lock. - sess, _ := concurrency.NewSession(cli.Client) + sess, _ := concurrency.NewSession(e.cli) defer sess.Close() muKey := path.Join(lockPrefix, index, field) @@ -1120,7 +960,7 @@ func (e *Etcd) RemoveShard(ctx context.Context, index, field string, shard uint6 } // Read shards again (within lock). - shards, err := e.shards(ctx, cli, index, field) + shards, err := e.shards(ctx, index, field) if err != nil { return errors.Wrap(err, "reading shards") } @@ -1137,7 +977,7 @@ func (e *Etcd) RemoveShard(ctx context.Context, index, field string, shard uint6 // If this is removing the last bit from the shards bitmap, then instead of // writing an empty bitmap, just delete the key. if shards.Count() == 0 { - _, err := cli.Delete(ctx, key) + _, err := e.cli.Delete(ctx, key) return err } @@ -1150,7 +990,7 @@ func (e *Etcd) RemoveShard(ctx context.Context, index, field string, shard uint6 op := clientv3.OpPut(key, "") op.WithValueBytes(buf.Bytes()) - if _, err := cli.Do(ctx, op); err != nil { + if _, err := e.cli.Do(ctx, op); err != nil { return errors.Wrap(err, "doing op") } From 0f4b273d3a2165dc6a90b7fe3374a2c053b79d9e Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Thu, 25 Feb 2021 11:30:09 +0100 Subject: [PATCH 2/8] Remove etcd cache --- etcd/cache.go | 94 ------------------------------------------------ etcd/embed.go | 16 +++++---- server/server.go | 2 +- 3 files changed, 10 insertions(+), 102 deletions(-) delete mode 100644 etcd/cache.go diff --git a/etcd/cache.go b/etcd/cache.go deleted file mode 100644 index 8005c87ef..000000000 --- a/etcd/cache.go +++ /dev/null @@ -1,94 +0,0 @@ -// Copyright 2017 Pilosa Corp. -// -// 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 etcd - -import ( - "context" - "sync" - "time" - - "github.com/pilosa/pilosa/v2/topology" -) - -// EtcdWithCache is a wrapper around the Etcd type which will return a -// cached value when the number of requests come in below a configured -// frequency. It also breaks the cache after a configured TTL. -type EtcdWithCache struct { - *Etcd - - peerMetadataMu sync.RWMutex - peerMetadata map[string][]byte - - peersMu sync.Mutex // peer-list cache updates - - nodes []*topology.Node // unmarshalled Node data - nodesTTL int // seconds - nodesLastRequest time.Time // last time requested -} - -// NewEtcdWithCache returns a new instance of Cache. -func NewEtcdWithCache(opt Options, replicas int) *EtcdWithCache { - return &EtcdWithCache{ - Etcd: NewEtcd(opt, replicas), - - nodesTTL: 6, - - peerMetadata: make(map[string][]byte), - } -} - -// Metadata is a cache wrapper around the Metadator.Metadata method. -func (c *EtcdWithCache) Metadata(ctx context.Context, peerID string) ([]byte, error) { - c.peerMetadataMu.RLock() - v, ok := c.peerMetadata[peerID] - c.peerMetadataMu.RUnlock() - if ok { - return v, nil - } - v, err := c.Etcd.Metadata(ctx, peerID) - if err == nil { - c.peerMetadataMu.Lock() - c.peerMetadata[peerID] = v - c.peerMetadataMu.Unlock() - } - return v, err -} - -// Nodes caches the result of the underlying implementation's node list. -func (c *EtcdWithCache) Nodes() []*topology.Node { - c.peersMu.Lock() - defer c.peersMu.Unlock() - - now := time.Now() - if now.Sub(c.nodesLastRequest) > (time.Duration(c.nodesTTL) * time.Second) { - c.nodes = c.Etcd.Nodes() - c.nodesLastRequest = now - } - return c.nodes -} - -// SetNodes implements the Noder interface as NOP -// (because we can't force to set nodes for etcd). -func (c *EtcdWithCache) SetNodes(nodes []*topology.Node) {} - -// AppendNode implements the Noder interface as NOP -// (because resizer is responsible for adding new nodes). -func (c *EtcdWithCache) AppendNode(node *topology.Node) {} - -// RemoveNode implements the Noder interface as NOP -// (because resizer is responsible for removing existing nodes) -func (c *EtcdWithCache) RemoveNode(nodeID string) bool { - return false -} diff --git a/etcd/embed.go b/etcd/embed.go index bef4689e2..cacac1cf7 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -106,6 +106,12 @@ func NewEtcd(opt Options, replicas int) *Etcd { func (e *Etcd) Close() error { _ = testhook.Closed(pilosa.NewAuditor(), e, nil) + if e.cli != nil { + if err := e.cli.Close(); err != nil { + log.Printf("Error closing etcd client: %v", err) + } + } + if e.e != nil { if e.resizeCancel != nil { e.resizeCancel() @@ -117,10 +123,6 @@ func (e *Etcd) Close() error { <-e.e.Server.StopNotify() } - if e.cli != nil { - return e.cli.Close() - } - return nil } @@ -250,15 +252,15 @@ func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, e return disco.NodeStateUnknown, err } - if len(resp.KVs) > 1 { + if len(resp.Kvs) > 1 { return disco.NodeStateUnknown, disco.ErrTooManyResults } - if len(resp.KVs) == 0 { + if len(resp.Kvs) == 0 { return disco.NodeStateUnknown, disco.ErrNoResults } - return disco.NodeState(resp.KVs[0].Value), nil + return disco.NodeState(resp.Kvs[0].Value), nil } func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, error) { diff --git a/server/server.go b/server/server.go index c07ba17a4..1317a8ebb 100644 --- a/server/server.go +++ b/server/server.go @@ -390,7 +390,7 @@ func (m *Command) SetupServer() error { m.Config.Etcd.Dir = filepath.Join(path, pilosa.DefaultDiscoDir) } - e := petcd.NewEtcdWithCache(m.Config.Etcd, m.Config.Cluster.ReplicaN) + e := petcd.NewEtcd(m.Config.Etcd, m.Config.Cluster.ReplicaN) discoOpt := pilosa.OptServerDisCo(e, e, e, e, e, e, e) serverOptions := []pilosa.ServerOption{ From 23f901635e0d82531eba5d278e47fbb52cdea426 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Thu, 25 Feb 2021 12:49:14 +0100 Subject: [PATCH 3/8] Revert "Remove etcd cache" This reverts commit 0f4b273d3a2165dc6a90b7fe3374a2c053b79d9e. --- etcd/cache.go | 94 ++++++++++++++++++++++++++++++++++++++++++++++++ etcd/embed.go | 25 +++++++------ server/server.go | 2 +- 3 files changed, 109 insertions(+), 12 deletions(-) create mode 100644 etcd/cache.go diff --git a/etcd/cache.go b/etcd/cache.go new file mode 100644 index 000000000..8005c87ef --- /dev/null +++ b/etcd/cache.go @@ -0,0 +1,94 @@ +// Copyright 2017 Pilosa Corp. +// +// 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 etcd + +import ( + "context" + "sync" + "time" + + "github.com/pilosa/pilosa/v2/topology" +) + +// EtcdWithCache is a wrapper around the Etcd type which will return a +// cached value when the number of requests come in below a configured +// frequency. It also breaks the cache after a configured TTL. +type EtcdWithCache struct { + *Etcd + + peerMetadataMu sync.RWMutex + peerMetadata map[string][]byte + + peersMu sync.Mutex // peer-list cache updates + + nodes []*topology.Node // unmarshalled Node data + nodesTTL int // seconds + nodesLastRequest time.Time // last time requested +} + +// NewEtcdWithCache returns a new instance of Cache. +func NewEtcdWithCache(opt Options, replicas int) *EtcdWithCache { + return &EtcdWithCache{ + Etcd: NewEtcd(opt, replicas), + + nodesTTL: 6, + + peerMetadata: make(map[string][]byte), + } +} + +// Metadata is a cache wrapper around the Metadator.Metadata method. +func (c *EtcdWithCache) Metadata(ctx context.Context, peerID string) ([]byte, error) { + c.peerMetadataMu.RLock() + v, ok := c.peerMetadata[peerID] + c.peerMetadataMu.RUnlock() + if ok { + return v, nil + } + v, err := c.Etcd.Metadata(ctx, peerID) + if err == nil { + c.peerMetadataMu.Lock() + c.peerMetadata[peerID] = v + c.peerMetadataMu.Unlock() + } + return v, err +} + +// Nodes caches the result of the underlying implementation's node list. +func (c *EtcdWithCache) Nodes() []*topology.Node { + c.peersMu.Lock() + defer c.peersMu.Unlock() + + now := time.Now() + if now.Sub(c.nodesLastRequest) > (time.Duration(c.nodesTTL) * time.Second) { + c.nodes = c.Etcd.Nodes() + c.nodesLastRequest = now + } + return c.nodes +} + +// SetNodes implements the Noder interface as NOP +// (because we can't force to set nodes for etcd). +func (c *EtcdWithCache) SetNodes(nodes []*topology.Node) {} + +// AppendNode implements the Noder interface as NOP +// (because resizer is responsible for adding new nodes). +func (c *EtcdWithCache) AppendNode(node *topology.Node) {} + +// RemoveNode implements the Noder interface as NOP +// (because resizer is responsible for removing existing nodes) +func (c *EtcdWithCache) RemoveNode(nodeID string) bool { + return false +} diff --git a/etcd/embed.go b/etcd/embed.go index cacac1cf7..53a96c9a0 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -37,6 +37,7 @@ import ( "go.etcd.io/etcd/clientv3/concurrency" "go.etcd.io/etcd/embed" "go.etcd.io/etcd/etcdserver/api/v3client" + "go.etcd.io/etcd/mvcc" "go.etcd.io/etcd/mvcc/mvccpb" "go.etcd.io/etcd/pkg/types" ) @@ -106,12 +107,6 @@ func NewEtcd(opt Options, replicas int) *Etcd { func (e *Etcd) Close() error { _ = testhook.Closed(pilosa.NewAuditor(), e, nil) - if e.cli != nil { - if err := e.cli.Close(); err != nil { - log.Printf("Error closing etcd client: %v", err) - } - } - if e.e != nil { if e.resizeCancel != nil { e.resizeCancel() @@ -123,6 +118,12 @@ func (e *Etcd) Close() error { <-e.e.Server.StopNotify() } + if e.cli != nil { + if err := e.cli.Close(); err != nil { + log.Printf("Error closing etcd client: %v", err) + } + } + return nil } @@ -239,7 +240,9 @@ func (e *Etcd) NodeState(ctx context.Context, peerID string) (disco.NodeState, e } func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, error) { - resp, err := e.cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly()) + kv := e.e.Server.KV() + + resp, err := kv.Range([]byte(path.Join(resizePrefix, peerID)), nil, mvcc.RangeOptions{Count: true}) if err != nil { return disco.NodeStateUnknown, err } @@ -247,20 +250,20 @@ func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, e return disco.NodeStateResizing, nil } - resp, err = e.cli.Get(ctx, path.Join(heartbeatPrefix, peerID)) + resp, err = kv.Range([]byte(path.Join(heartbeatPrefix, peerID)), nil, mvcc.RangeOptions{}) if err != nil { return disco.NodeStateUnknown, err } - if len(resp.Kvs) > 1 { + if len(resp.KVs) > 1 { return disco.NodeStateUnknown, disco.ErrTooManyResults } - if len(resp.Kvs) == 0 { + if len(resp.KVs) == 0 { return disco.NodeStateUnknown, disco.ErrNoResults } - return disco.NodeState(resp.Kvs[0].Value), nil + return disco.NodeState(resp.KVs[0].Value), nil } func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, error) { diff --git a/server/server.go b/server/server.go index 1317a8ebb..c07ba17a4 100644 --- a/server/server.go +++ b/server/server.go @@ -390,7 +390,7 @@ func (m *Command) SetupServer() error { m.Config.Etcd.Dir = filepath.Join(path, pilosa.DefaultDiscoDir) } - e := petcd.NewEtcd(m.Config.Etcd, m.Config.Cluster.ReplicaN) + e := petcd.NewEtcdWithCache(m.Config.Etcd, m.Config.Cluster.ReplicaN) discoOpt := pilosa.OptServerDisCo(e, e, e, e, e, e, e) serverOptions := []pilosa.ServerOption{ From 1623007af18e95bbefdfca65e5328abb0b4d8021 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Thu, 25 Feb 2021 13:39:47 +0100 Subject: [PATCH 4/8] Add waitgroup - don't close the server wait for all keepaliveFunc --- etcd/embed.go | 26 ++++++++++++++++---------- 1 file changed, 16 insertions(+), 10 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index 53a96c9a0..c899895d2 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -24,6 +24,7 @@ import ( "path" "sort" "strings" + "sync" "time" "github.com/pilosa/pilosa/v2" @@ -37,7 +38,6 @@ import ( "go.etcd.io/etcd/clientv3/concurrency" "go.etcd.io/etcd/embed" "go.etcd.io/etcd/etcdserver/api/v3client" - "go.etcd.io/etcd/mvcc" "go.etcd.io/etcd/mvcc/mvccpb" "go.etcd.io/etcd/pkg/types" ) @@ -93,12 +93,14 @@ type Etcd struct { e *embed.Etcd cli *clientv3.Client + wg *sync.WaitGroup } func NewEtcd(opt Options, replicas int) *Etcd { e := &Etcd{ options: opt, replicas: replicas, + wg: &sync.WaitGroup{}, } return e } @@ -114,6 +116,8 @@ func (e *Etcd) Close() error { if e.heartbeatCancel != nil { e.heartbeatCancel() } + + e.wg.Wait() e.e.Close() <-e.e.Server.StopNotify() } @@ -240,9 +244,7 @@ func (e *Etcd) NodeState(ctx context.Context, peerID string) (disco.NodeState, e } func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, error) { - kv := e.e.Server.KV() - - resp, err := kv.Range([]byte(path.Join(resizePrefix, peerID)), nil, mvcc.RangeOptions{Count: true}) + resp, err := e.cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly()) if err != nil { return disco.NodeStateUnknown, err } @@ -250,20 +252,20 @@ func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, e return disco.NodeStateResizing, nil } - resp, err = kv.Range([]byte(path.Join(heartbeatPrefix, peerID)), nil, mvcc.RangeOptions{}) + resp, err = e.cli.Get(ctx, path.Join(heartbeatPrefix, peerID)) if err != nil { return disco.NodeStateUnknown, err } - if len(resp.KVs) > 1 { + if len(resp.Kvs) > 1 { return disco.NodeStateUnknown, disco.ErrTooManyResults } - if len(resp.KVs) == 0 { + if len(resp.Kvs) == 0 { return disco.NodeStateUnknown, disco.ErrNoResults } - return disco.NodeState(resp.KVs[0].Value), nil + return disco.NodeState(resp.Kvs[0].Value), nil } func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, error) { @@ -723,7 +725,10 @@ func (e *Etcd) leaseKeepAlive(ctx context.Context, ttl int64) (clientv3.LeaseID, keepaliveFunc := func(tick time.Duration) { ticker := time.NewTicker(tick) - defer ticker.Stop() + defer func() { + ticker.Stop() + e.wg.Done() + }() for { select { @@ -733,7 +738,6 @@ func (e *Etcd) leaseKeepAlive(ctx context.Context, ttl int64) (clientv3.LeaseID, // here, resulting in massive piles of excess goroutines. revoker, cancel := context.WithTimeout(context.Background(), time.Duration(ttl)) defer cancel() - if _, err := e.cli.Revoke(revoker, leaseResp.ID); err != nil { log.Printf("leaseKeepAlive: revokes the lease (ID: %x): %#v\n", leaseResp.ID, err) } @@ -745,6 +749,8 @@ func (e *Etcd) leaseKeepAlive(ctx context.Context, ttl int64) (clientv3.LeaseID, } } } + + e.wg.Add(1) go keepaliveFunc(time.Second) return leaseResp.ID, ctx, cancelFunc, nil From a5f3bce3bffad2ed9632c717146f7968726351a4 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Thu, 25 Feb 2021 16:56:30 +0100 Subject: [PATCH 5/8] Use hookedClient --- etcd/embed.go | 15 +++++++-------- 1 file changed, 7 insertions(+), 8 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index c899895d2..a8c04b5d1 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -92,7 +92,7 @@ type Etcd struct { lm leaseMetadata e *embed.Etcd - cli *clientv3.Client + cli *hookedClient wg *sync.WaitGroup } @@ -123,9 +123,7 @@ func (e *Etcd) Close() error { } if e.cli != nil { - if err := e.cli.Close(); err != nil { - log.Printf("Error closing etcd client: %v", err) - } + e.cli.Close() } return nil @@ -203,7 +201,7 @@ func (e *Etcd) Start(ctx context.Context) (disco.InitialClusterState, error) { } _ = testhook.Opened(pilosa.NewAuditor(), e, nil) e.e = etcd - e.cli = v3client.New(e.e.Server) + e.cli = &hookedClient{Client: v3client.New(e.e.Server)} select { case <-ctx.Done(): @@ -738,6 +736,7 @@ func (e *Etcd) leaseKeepAlive(ctx context.Context, ttl int64) (clientv3.LeaseID, // here, resulting in massive piles of excess goroutines. revoker, cancel := context.WithTimeout(context.Background(), time.Duration(ttl)) defer cancel() + if _, err := e.cli.Revoke(revoker, leaseResp.ID); err != nil { log.Printf("leaseKeepAlive: revokes the lease (ID: %x): %#v\n", leaseResp.ID, err) } @@ -837,7 +836,7 @@ func (e *Etcd) AddShards(ctx context.Context, index, field string, shards *roari // } // Create a session to acquire a lock. - sess, _ := concurrency.NewSession(e.cli) + sess, _ := concurrency.NewSession(e.cli.Client) defer sess.Close() muKey := path.Join(lockPrefix, index, field) @@ -896,7 +895,7 @@ func (e *Etcd) AddShard(ctx context.Context, index, field string, shard uint64) // write shards to etcd. // Create a session to acquire a lock. - sess, _ := concurrency.NewSession(e.cli) + sess, _ := concurrency.NewSession(e.cli.Client) defer sess.Close() muKey := path.Join(lockPrefix, index, field) @@ -959,7 +958,7 @@ func (e *Etcd) RemoveShard(ctx context.Context, index, field string, shard uint6 // write shards to etcd. // Create a session to acquire a lock. - sess, _ := concurrency.NewSession(e.cli) + sess, _ := concurrency.NewSession(e.cli.Client) defer sess.Close() muKey := path.Join(lockPrefix, index, field) From aa15e0855838ab2a256b333f64955942a9bbb9d8 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Thu, 25 Feb 2021 18:12:13 +0100 Subject: [PATCH 6/8] Reduce number of Txn --- etcd/embed.go | 39 ++++++++++++++------------------------- 1 file changed, 14 insertions(+), 25 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index a8c04b5d1..70a4b14c6 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -18,7 +18,6 @@ import ( "bytes" "context" "encoding/json" - "fmt" "log" "net" "path" @@ -243,6 +242,8 @@ func (e *Etcd) NodeState(ctx context.Context, peerID string) (disco.NodeState, e func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, error) { resp, err := e.cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly()) + // kv := e.e.Server.KV() + // resp, err := kv.Range([]byte(path.Join(resizePrefix, peerID)), nil, mvcc.RangeOptions{Count: true}) if err != nil { return disco.NodeStateUnknown, err } @@ -251,19 +252,21 @@ func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, e } resp, err = e.cli.Get(ctx, path.Join(heartbeatPrefix, peerID)) + // resp, err = kv.Range([]byte(path.Join(heartbeatPrefix, peerID)), nil, mvcc.RangeOptions{}) if err != nil { return disco.NodeStateUnknown, err } + kvs := resp.Kvs - if len(resp.Kvs) > 1 { + if len(kvs) > 1 { return disco.NodeStateUnknown, disco.ErrTooManyResults } - if len(resp.Kvs) == 0 { + if len(kvs) == 0 { return disco.NodeStateUnknown, disco.ErrNoResults } - return disco.NodeState(resp.Kvs[0].Value), nil + return disco.NodeState(kvs[0].Value), nil } func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, error) { @@ -658,28 +661,19 @@ func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { } func (e *Etcd) getKey(ctx context.Context, key string) ([]string, [][]byte, error) { - resp, err := e.cli.Txn(ctx). - If(clientv3.Compare(clientv3.Version(key), ">", -1)). - Then(clientv3.OpGet(key, clientv3.WithPrefix())). - Commit() + resp, err := e.cli.Get(ctx, key, clientv3.WithPrefix()) if err != nil { return nil, nil, err } - if !resp.Succeeded { - return nil, nil, fmt.Errorf("key %s does not exist", key) - } - var ( keys []string values [][]byte ) - for _, r := range resp.Responses { - for _, kv := range r.GetResponseRange().Kvs { - keys = append(keys, string(kv.Key)) - values = append(values, kv.Value) - } + for _, kv := range resp.Kvs { + keys = append(keys, string(kv.Key)) + values = append(values, kv.Value) } return keys, values, nil @@ -697,16 +691,11 @@ func (e *Etcd) keyExists(ctx context.Context, key string) (bool, error) { } func (e *Etcd) delKey(ctx context.Context, key string, withPrefix bool) (err error) { - var opts []clientv3.OpOption if withPrefix { - opts = append(opts, clientv3.WithPrefix()) + _, err = e.cli.Delete(ctx, key, clientv3.WithPrefix()) + } else { + _, err = e.cli.Delete(ctx, key) } - - _, err = e.cli.KV.Txn(ctx). - If(clientv3.Compare(clientv3.Version(key), ">", -1)). - Then(clientv3.OpDelete(key, opts...)). - Commit() - return err } From 49dc48f0574f534ac42d7b7b1a3c8e1a8b387927 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Thu, 25 Feb 2021 18:34:24 +0100 Subject: [PATCH 7/8] Switch to server API for KV Get/Range --- etcd/embed.go | 60 +++++++++++++++++++++++++++++++++++++-------------- 1 file changed, 44 insertions(+), 16 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index 70a4b14c6..1079bf6bb 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -37,6 +37,7 @@ import ( "go.etcd.io/etcd/clientv3/concurrency" "go.etcd.io/etcd/embed" "go.etcd.io/etcd/etcdserver/api/v3client" + "go.etcd.io/etcd/mvcc" "go.etcd.io/etcd/mvcc/mvccpb" "go.etcd.io/etcd/pkg/types" ) @@ -241,9 +242,9 @@ func (e *Etcd) NodeState(ctx context.Context, peerID string) (disco.NodeState, e } func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, error) { - resp, err := e.cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly()) - // kv := e.e.Server.KV() - // resp, err := kv.Range([]byte(path.Join(resizePrefix, peerID)), nil, mvcc.RangeOptions{Count: true}) + // resp, err := e.cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly()) + kv := e.e.Server.KV() + resp, err := kv.Range([]byte(path.Join(resizePrefix, peerID)), nil, mvcc.RangeOptions{Count: true}) if err != nil { return disco.NodeStateUnknown, err } @@ -251,12 +252,12 @@ func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, e return disco.NodeStateResizing, nil } - resp, err = e.cli.Get(ctx, path.Join(heartbeatPrefix, peerID)) - // resp, err = kv.Range([]byte(path.Join(heartbeatPrefix, peerID)), nil, mvcc.RangeOptions{}) + // resp, err = e.cli.Get(ctx, path.Join(heartbeatPrefix, peerID)) + resp, err = kv.Range([]byte(path.Join(heartbeatPrefix, peerID)), nil, mvcc.RangeOptions{}) if err != nil { return disco.NodeStateUnknown, err } - kvs := resp.Kvs + kvs := resp.KVs if len(kvs) > 1 { return disco.NodeStateUnknown, disco.ErrTooManyResults @@ -501,20 +502,23 @@ func (e *Etcd) Schema(ctx context.Context) (disco.Schema, error) { } func (e *Etcd) Metadata(ctx context.Context, peerID string) ([]byte, error) { - resp, err := e.cli.Get(ctx, path.Join(metadataPrefix, peerID)) + // resp, err := e.cli.Get(ctx, path.Join(metadataPrefix, peerID)) + kv := e.e.Server.KV() + resp, err := kv.Range([]byte(path.Join(metadataPrefix, peerID)), nil, mvcc.RangeOptions{}) if err != nil { return nil, err } + kvs := resp.KVs - if len(resp.Kvs) > 1 { + if len(kvs) > 1 { return nil, disco.ErrTooManyResults } - if len(resp.Kvs) == 0 { + if len(kvs) == 0 { return nil, disco.ErrNoResults } - return resp.Kvs[0].Value, nil + return kvs[0].Value, nil } func (e *Etcd) SetMetadata(ctx context.Context, metadata []byte) error { @@ -647,31 +651,53 @@ func (e *Etcd) putKey(ctx context.Context, key, val string, opts ...clientv3.OpO func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { // Get the current value for the key. - resp, err := e.cli.Get(ctx, key) + // resp, err := e.cli.Get(ctx, key) + kv := e.e.Server.KV() + resp, err := kv.Range([]byte(key), nil, mvcc.RangeOptions{}) if err != nil { return nil, err } + kvs := resp.KVs // TODO: consider returning a "key does not exist" error instead of (nil, nil) - if len(resp.Kvs) == 0 { + if len(kvs) == 0 { return nil, nil } - return resp.Kvs[0].Value, nil + return kvs[0].Value, nil } func (e *Etcd) getKey(ctx context.Context, key string) ([]string, [][]byte, error) { - resp, err := e.cli.Get(ctx, key, clientv3.WithPrefix()) + getPrefix := func(key []byte) []byte { + end := make([]byte, len(key)) + copy(end, key) + for i := len(end) - 1; i >= 0; i-- { + if end[i] < 0xff { + end[i] = end[i] + 1 + end = end[:i+1] + return end + } + } + // next prefix does not exist (e.g., 0xffff); + // default to WithFromKey policy + return nil + } + + kv := e.e.Server.KV() + resp, err := kv.Range([]byte(key), getPrefix([]byte(key)), mvcc.RangeOptions{}) + + // resp, err := e.cli.Get(ctx, key, clientv3.WithPrefix()) if err != nil { return nil, nil, err } + kvs := resp.KVs var ( keys []string values [][]byte ) - for _, kv := range resp.Kvs { + for _, kv := range kvs { keys = append(keys, string(kv.Key)) values = append(values, kv.Value) } @@ -680,7 +706,9 @@ func (e *Etcd) getKey(ctx context.Context, key string) ([]string, [][]byte, erro } func (e *Etcd) keyExists(ctx context.Context, key string) (bool, error) { - resp, err := e.cli.Get(ctx, key, clientv3.WithCountOnly()) + // resp, err := e.cli.Get(ctx, key, clientv3.WithCountOnly()) + kv := e.e.Server.KV() + resp, err := kv.Range([]byte(key), nil, mvcc.RangeOptions{Count: true}) if err != nil { return false, err } From 50cbb72619271fcf4d22f4ea9087a937e3035b94 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Thu, 25 Feb 2021 19:00:52 +0100 Subject: [PATCH 8/8] Remove comments/leftovers --- etcd/embed.go | 31 ++++--------------------------- 1 file changed, 4 insertions(+), 27 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index 1079bf6bb..06c081e74 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -242,7 +242,6 @@ func (e *Etcd) NodeState(ctx context.Context, peerID string) (disco.NodeState, e } func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, error) { - // resp, err := e.cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly()) kv := e.e.Server.KV() resp, err := kv.Range([]byte(path.Join(resizePrefix, peerID)), nil, mvcc.RangeOptions{Count: true}) if err != nil { @@ -252,7 +251,6 @@ func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, e return disco.NodeStateResizing, nil } - // resp, err = e.cli.Get(ctx, path.Join(heartbeatPrefix, peerID)) resp, err = kv.Range([]byte(path.Join(heartbeatPrefix, peerID)), nil, mvcc.RangeOptions{}) if err != nil { return disco.NodeStateUnknown, err @@ -447,7 +445,7 @@ func (e *Etcd) DeleteNode(ctx context.Context, nodeID string) error { } func (e *Etcd) Schema(ctx context.Context) (disco.Schema, error) { - keys, vals, err := e.getKey(ctx, schemaPrefix) + keys, vals, err := e.getKeyWithPrefix(ctx, schemaPrefix) if err != nil { return nil, err } @@ -502,7 +500,6 @@ func (e *Etcd) Schema(ctx context.Context) (disco.Schema, error) { } func (e *Etcd) Metadata(ctx context.Context, peerID string) ([]byte, error) { - // resp, err := e.cli.Get(ctx, path.Join(metadataPrefix, peerID)) kv := e.e.Server.KV() resp, err := kv.Range([]byte(path.Join(metadataPrefix, peerID)), nil, mvcc.RangeOptions{}) if err != nil { @@ -651,7 +648,6 @@ func (e *Etcd) putKey(ctx context.Context, key, val string, opts ...clientv3.OpO func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { // Get the current value for the key. - // resp, err := e.cli.Get(ctx, key) kv := e.e.Server.KV() resp, err := kv.Range([]byte(key), nil, mvcc.RangeOptions{}) if err != nil { @@ -667,30 +663,12 @@ func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { return kvs[0].Value, nil } -func (e *Etcd) getKey(ctx context.Context, key string) ([]string, [][]byte, error) { - getPrefix := func(key []byte) []byte { - end := make([]byte, len(key)) - copy(end, key) - for i := len(end) - 1; i >= 0; i-- { - if end[i] < 0xff { - end[i] = end[i] + 1 - end = end[:i+1] - return end - } - } - // next prefix does not exist (e.g., 0xffff); - // default to WithFromKey policy - return nil - } - - kv := e.e.Server.KV() - resp, err := kv.Range([]byte(key), getPrefix([]byte(key)), mvcc.RangeOptions{}) - - // resp, err := e.cli.Get(ctx, key, clientv3.WithPrefix()) +func (e *Etcd) getKeyWithPrefix(ctx context.Context, key string) ([]string, [][]byte, error) { + resp, err := e.cli.Get(ctx, key, clientv3.WithPrefix()) if err != nil { return nil, nil, err } - kvs := resp.KVs + kvs := resp.Kvs var ( keys []string @@ -706,7 +684,6 @@ func (e *Etcd) getKey(ctx context.Context, key string) ([]string, [][]byte, erro } func (e *Etcd) keyExists(ctx context.Context, key string) (bool, error) { - // resp, err := e.cli.Get(ctx, key, clientv3.WithCountOnly()) kv := e.e.Server.KV() resp, err := kv.Range([]byte(key), nil, mvcc.RangeOptions{Count: true}) if err != nil {