One shared etcd client

This commit is contained in:
Kuba Podgórski 2021-02-24 22:05:49 +01:00
parent bffdee1e8e
commit 01e0c44069

View file

@ -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")
}