From 7a1e4cd2d7076403287c18fa674310f617d2e06b Mon Sep 17 00:00:00 2001 From: Antonio Navarro Perez Date: Thu, 4 Mar 2021 12:28:01 +0100 Subject: [PATCH] Improve leased keys. Added a leasedKV struct in charge of maintain a lease for a specific key. Signed-off-by: Antonio Navarro Perez --- etcd/embed.go | 191 ++++++----------------------------------------- etcd/leasedkv.go | 119 +++++++++++++++++++++++++++++ 2 files changed, 143 insertions(+), 167 deletions(-) create mode 100644 etcd/leasedkv.go diff --git a/etcd/embed.go b/etcd/embed.go index ee0fd1e55..0ad0d7b6b 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -23,8 +23,6 @@ import ( "path" "sort" "strings" - "sync" - "time" "github.com/pilosa/pilosa/v2/disco" "github.com/pilosa/pilosa/v2/roaring" @@ -35,7 +33,6 @@ import ( "go.etcd.io/etcd/clientv3/concurrency" "go.etcd.io/etcd/embed" "go.etcd.io/etcd/etcdserver/api/v3client" - "go.etcd.io/etcd/etcdserver/api/v3rpc/rpctypes" "go.etcd.io/etcd/mvcc/mvccpb" "go.etcd.io/etcd/pkg/types" ) @@ -74,31 +71,20 @@ const ( lockPrefix = "/lock/" ) -type leaseMetadata struct { - started bool -} - type Etcd struct { options Options replicas int - heartbeatID clientv3.LeaseID - heartbeatCancel context.CancelFunc - - resizeCancel context.CancelFunc - - lm leaseMetadata - e *embed.Etcd cli *clientv3.Client - wg *sync.WaitGroup + + heartBeatLeasedKV, resizeLeasedKV *leasedKV } func NewEtcd(opt Options, replicas int) *Etcd { e := &Etcd{ options: opt, replicas: replicas, - wg: &sync.WaitGroup{}, } if e.options.HeartbeatTTL == 0 { @@ -110,14 +96,14 @@ func NewEtcd(opt Options, replicas int) *Etcd { // Close implements io.Closer func (e *Etcd) Close() error { if e.e != nil { - if e.resizeCancel != nil { - e.resizeCancel() + if e.resizeLeasedKV != nil { + e.resizeLeasedKV.Stop() + e.resizeLeasedKV = nil } - if e.heartbeatCancel != nil { - e.heartbeatCancel() + if e.heartBeatLeasedKV != nil { + e.heartBeatLeasedKV.Stop() } - e.wg.Wait() e.e.Close() <-e.e.Server.StopNotify() } @@ -210,38 +196,16 @@ func (e *Etcd) Start(ctx context.Context) (disco.InitialClusterState, error) { return state, err case <-e.e.Server.ReadyNotify(): - return state, e.startHeartbeat() + return state, e.startHeartbeat(ctx) } } -func (e *Etcd) startHeartbeat() error { - ctx, heartbeatCancel := context.WithCancel(context.Background()) - e.heartbeatCancel = heartbeatCancel +func (e *Etcd) startHeartbeat(ctx context.Context) error { + key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.NodeStateStarting + e.heartBeatLeasedKV = newLeasedKV(e.cli, key, e.options.HeartbeatTTL) - cb := func(heartbeatID clientv3.LeaseID) error { - key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.NodeStateStarting - if e.e.Config().ClusterState == embed.ClusterStateFlagExisting { - value = disco.NodeStateResizing - } else if e.lm.started { - value = disco.NodeStateStarted - } - - if _, err := e.cli.Txn(ctx). - Then(clientv3.OpPut(key, string(value), clientv3.WithLease(heartbeatID))). - Commit(); err != nil { - - heartbeatCancel() - return errors.Wrapf(err, "startHeartbeat: txn puts a key-value (%s, %s) with lease (%v)", key, value, heartbeatID) - } - - e.heartbeatID = heartbeatID - - return nil - } - - _, err := e.leaseKeepAlive(ctx, heartbeatCancel, cb) - if err != nil { - return errors.Wrap(err, "startHeartbeat: creates a new heartbeat") + if err := e.heartBeatLeasedKV.Start(ctx, string(value)); err != nil { + return errors.Wrap(err, "startHeartbeat: starting a new heartbeat") } return nil @@ -319,15 +283,7 @@ func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, erro } func (e *Etcd) Started(ctx context.Context) (err error) { - key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.NodeStateStarted - _, err = e.cli.Txn(ctx). - Then(clientv3.OpPut(key, string(value), clientv3.WithLease(e.heartbeatID))). - Commit() - - if err == nil { - e.lm.started = true - } - return err + return e.heartBeatLeasedKV.Set(ctx, string(disco.NodeStateStarted)) } func (e *Etcd) ID() string { @@ -411,43 +367,27 @@ func (e *Etcd) ClusterState(ctx context.Context) (disco.ClusterState, error) { } func (e *Etcd) Resize(ctx context.Context) (func([]byte) error, error) { - ctx, resizeCancel := context.WithCancel(ctx) + key := path.Join(resizePrefix, e.e.Server.ID().String()) + if e.resizeLeasedKV == nil { + e.resizeLeasedKV = newLeasedKV(e.cli, key, e.options.HeartbeatTTL) + } - cb := func(clientv3.LeaseID) error { return nil } - - resizeID, err := e.leaseKeepAlive(ctx, resizeCancel, cb) - if err != nil { + if err := e.resizeLeasedKV.Start(ctx, ""); err != nil { return nil, errors.Wrap(err, "Resize: creates a new hearbeat") } - // Check if key exists - maybe we are still resizing - key := path.Join(resizePrefix, e.e.Server.ID().String()) - txnResp, err := e.cli.Txn(ctx). - If(clientv3util.KeyMissing(key)). - Then(clientv3.OpPut(key, "", clientv3.WithLease(resizeID))). - Commit() - if err != nil { - resizeCancel() - return nil, errors.Wrapf(err, "Resize: txn puts key (%s) with lease (%v)", key, resizeID) - } - - if !txnResp.Succeeded { - resizeCancel() - return nil, errors.Errorf("Resize: key (%s) exists - maybe node (%s) is resizing", key, e.ID()) - } - - e.resizeCancel = resizeCancel - return func(value []byte) error { log.Println("Update progress:", key, string(value)) - return e.putKey(ctx, key, string(value), clientv3.WithLease(resizeID)) + return e.putKey(ctx, key, string(value), clientv3.WithIgnoreLease()) }, nil } func (e *Etcd) DoneResize() error { - if e.resizeCancel != nil { - e.resizeCancel() + if e.resizeLeasedKV != nil { + e.resizeLeasedKV.Stop() } + + e.resizeLeasedKV = nil return nil } @@ -768,89 +708,6 @@ func (e *Etcd) delKey(ctx context.Context, key string, withPrefix bool) (err err return err } -// leaseKeepAlive creates a lease with the given ttl (treated as a time.Duration), -// 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, cancelFunc context.CancelFunc, cb func(clientv3.LeaseID) error) (clientv3.LeaseID, error) { - leaseResp, err := e.cli.Grant(ctx, e.options.HeartbeatTTL) - if err != nil { - cancelFunc() - return 0, errors.Wrapf(err, "leaseKeepAlive: creates a new lease (TTL: %d s.)", e.options.HeartbeatTTL) - } - - keepaliveFunc := func(tick time.Duration) error { - ticker := time.NewTicker(tick) - defer func() { - ticker.Stop() - e.wg.Done() - }() - - // leaseResp is a var within the function because we may need to reset - // it later if the lease has to be re-granted. - var leaseResp *clientv3.LeaseGrantResponse = leaseResp - - for { - select { - case <-ctx.Done(): - // Because of the load balancer, this can take ridiculously - // long times to run if the cluster's already down when we get - // here, resulting in massive piles of excess goroutines. - revoker, cancel := context.WithTimeout(context.Background(), time.Duration(e.options.HeartbeatTTL)*time.Second) - 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) - return errors.Wrap(err, "revoking lease") - } - return nil - case <-ticker.C: - _, err := e.cli.KeepAliveOnce(ctx, leaseResp.ID) - if err == rpctypes.ErrLeaseNotFound { - // We create a new client here because in the case where we - // have lost track of the lease, it's likely that we've also - // lost the client at e.cli. - // TODO: should this close/reset e.cli instead? - cli := v3client.New(e.e.Server) - var err error - leaseResp, err = cli.Grant(ctx, e.options.HeartbeatTTL) - cli.Close() - if err != nil { - cancelFunc() - return errors.Wrapf(err, "leaseKeepAlive: creates a new lease (TTL: %d s.)", e.options.HeartbeatTTL) - } - - // Call the callback. - if err := cb(leaseResp.ID); err != nil { - cancelFunc() - return errors.Wrap(err, "calling callback") - } - - // TODO: this can't be here in this general function because resize doesn't need this. - if err := e.Started(ctx); err != nil { - cancelFunc() - return errors.Wrap(err, "setting to started") - } - } else if err != nil { - log.Printf("leaseKeepAlive: renews the lease (ID: %x): %v\n", leaseResp.ID, err) - } - } - } - } - - e.wg.Add(1) - go func() { - if err := keepaliveFunc(time.Second); err != nil { - log.Printf("leaseKeepAlive: goroutine err: %v\n", err) - } - }() - - if err := cb(leaseResp.ID); err != nil { - return 0, errors.Wrap(err, "calling callback") - } - - return leaseResp.ID, nil -} - func memberList(cli *clientv3.Client) (ids []uint64, names []string, urls []string) { ml, err := cli.MemberList(context.TODO()) if err != nil { diff --git a/etcd/leasedkv.go b/etcd/leasedkv.go new file mode 100644 index 000000000..6f0082310 --- /dev/null +++ b/etcd/leasedkv.go @@ -0,0 +1,119 @@ +// 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" + + "github.com/pilosa/pilosa/v2/disco" + "github.com/pkg/errors" + "go.etcd.io/etcd/clientv3" + "go.etcd.io/etcd/clientv3/clientv3util" +) + +type leasedKV struct { + cli *clientv3.Client + cancel context.CancelFunc + + mu sync.Mutex + + key, value string + ttlSeconds int64 +} + +func newLeasedKV(cli *clientv3.Client, key string, ttlSeconds int64) *leasedKV { + return &leasedKV{ + cli: cli, + key: key, + ttlSeconds: ttlSeconds, + } +} + +func (l *leasedKV) Start(ctx context.Context, initValue string) error { + l.mu.Lock() + defer l.mu.Unlock() + + ctx, cancel := context.WithCancel(ctx) + l.cancel = cancel + + return l.create(ctx, initValue) +} + +func (l *leasedKV) create(ctx context.Context, initValue string) error { + leaseResp, err := l.cli.Grant(ctx, l.ttlSeconds) + if err != nil { + return errors.Wrap(err, "creating a lease") + } + + if _, err := l.cli.Txn(ctx). + Then(clientv3.OpPut(l.key, initValue, clientv3.WithLease(leaseResp.ID))). + Commit(); err != nil { + return errors.Wrapf(err, "creating key %s with value [%s]", l.key, initValue) + } + + if _, err := l.cli.KeepAlive(ctx, leaseResp.ID); err != nil { + return errors.Wrapf(err, "keeping alive the lease for the key %s with value %s", l.key, l.value) + } + + l.value = initValue + + return nil +} + +func (l *leasedKV) Stop() { + l.mu.Lock() + defer l.mu.Unlock() + + if l.cancel != nil { + l.cancel() + } +} + +func (l *leasedKV) Set(ctx context.Context, value string) error { + l.mu.Lock() + defer l.mu.Unlock() + + if _, err := l.cli.Txn(ctx). + Then(clientv3.OpPut(l.key, value, clientv3.WithIgnoreLease())). + Commit(); err != nil { + return errors.Wrapf(err, "creating key %s with value [%s]", l.key, l.value) + } + + l.value = value + + return nil +} + +func (l *leasedKV) Get(ctx context.Context) (string, error) { + l.mu.Lock() + defer l.mu.Unlock() + + getResp, err := l.cli.Txn(ctx). + If(clientv3util.KeyExists(l.key)). + Then(clientv3.OpGet(l.key, clientv3.WithIgnoreLease())). + Commit() + if err != nil { + return "", errors.Wrapf(err, "getting key %s", l.key) + } + + if !getResp.Succeeded || len(getResp.Responses) == 0 { + return "", disco.ErrNoResults + } + + l.value = string(getResp.Responses[0].GetResponseRange().Kvs[0].Value) + + return l.value, nil +}