From 7a1e4cd2d7076403287c18fa674310f617d2e06b Mon Sep 17 00:00:00 2001 From: Antonio Navarro Perez Date: Thu, 4 Mar 2021 12:28:01 +0100 Subject: [PATCH 1/8] 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 +} From 851c41cd45d5825d08dcb130e09f55a0e58e9245 Mon Sep 17 00:00:00 2001 From: Antonio Navarro Perez Date: Fri, 5 Mar 2021 13:49:18 +0100 Subject: [PATCH 2/8] Implement infinite lease renewal logic. Signed-off-by: Antonio Navarro Perez --- etcd/embed.go | 6 +-- etcd/leasedkv.go | 61 ++++++++++++++++++++++++++---- etcd/leasedkv_test.go | 87 +++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 144 insertions(+), 10 deletions(-) create mode 100644 etcd/leasedkv_test.go diff --git a/etcd/embed.go b/etcd/embed.go index 0ad0d7b6b..db8e5672f 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -201,10 +201,10 @@ func (e *Etcd) Start(ctx context.Context) (disco.InitialClusterState, error) { } func (e *Etcd) startHeartbeat(ctx context.Context) error { - key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.NodeStateStarting + key := heartbeatPrefix + e.e.Server.ID().String() e.heartBeatLeasedKV = newLeasedKV(e.cli, key, e.options.HeartbeatTTL) - if err := e.heartBeatLeasedKV.Start(ctx, string(value)); err != nil { + if err := e.heartBeatLeasedKV.Start(string(disco.NodeStateStarting)); err != nil { return errors.Wrap(err, "startHeartbeat: starting a new heartbeat") } @@ -372,7 +372,7 @@ func (e *Etcd) Resize(ctx context.Context) (func([]byte) error, error) { e.resizeLeasedKV = newLeasedKV(e.cli, key, e.options.HeartbeatTTL) } - if err := e.resizeLeasedKV.Start(ctx, ""); err != nil { + if err := e.resizeLeasedKV.Start(""); err != nil { return nil, errors.Wrap(err, "Resize: creates a new hearbeat") } diff --git a/etcd/leasedkv.go b/etcd/leasedkv.go index 6f0082310..77df1c73c 100644 --- a/etcd/leasedkv.go +++ b/etcd/leasedkv.go @@ -16,7 +16,9 @@ package etcd import ( "context" + "log" "sync" + "time" "github.com/pilosa/pilosa/v2/disco" "github.com/pkg/errors" @@ -31,6 +33,7 @@ type leasedKV struct { mu sync.Mutex key, value string + stopped bool ttlSeconds int64 } @@ -42,17 +45,17 @@ func newLeasedKV(cli *clientv3.Client, key string, ttlSeconds int64) *leasedKV { } } -func (l *leasedKV) Start(ctx context.Context, initValue string) error { +func (l *leasedKV) Start(initValue string) error { l.mu.Lock() defer l.mu.Unlock() - ctx, cancel := context.WithCancel(ctx) - l.cancel = cancel - - return l.create(ctx, initValue) + return l.create(initValue) } -func (l *leasedKV) create(ctx context.Context, initValue string) error { +func (l *leasedKV) create(initValue string) error { + ctx, cancel := context.WithCancel(context.Background()) + l.cancel = cancel + leaseResp, err := l.cli.Grant(ctx, l.ttlSeconds) if err != nil { return errors.Wrap(err, "creating a lease") @@ -64,15 +67,44 @@ func (l *leasedKV) create(ctx context.Context, initValue string) error { return errors.Wrapf(err, "creating key %s with value [%s]", l.key, initValue) } - if _, err := l.cli.KeepAlive(ctx, leaseResp.ID); err != nil { + kaChann, err := l.cli.KeepAlive(ctx, leaseResp.ID) + if err != nil { return errors.Wrapf(err, "keeping alive the lease for the key %s with value %s", l.key, l.value) } l.value = initValue + go l.consumeLease(kaChann) + return nil } +func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) { + for { + kresp, ok := <-ch + if !ok { + if l.stopped { + return + } + + l.mu.Lock() + defer l.mu.Unlock() + + retry(1*time.Second, func() error { + return l.create(l.value) + }) + + log.Println("lease recreated after a problem. Key:", l.key) + + return + } + + if kresp == nil { + continue + } + } +} + func (l *leasedKV) Stop() { l.mu.Lock() defer l.mu.Unlock() @@ -80,6 +112,8 @@ func (l *leasedKV) Stop() { if l.cancel != nil { l.cancel() } + + l.stopped = true } func (l *leasedKV) Set(ctx context.Context, value string) error { @@ -117,3 +151,16 @@ func (l *leasedKV) Get(ctx context.Context) (string, error) { return l.value, nil } + +func retry(sleep time.Duration, f func() error) { + for { + err := f() + if err == nil { + return + } + + time.Sleep(sleep) + + log.Println("retrying after error:", err) + } +} diff --git a/etcd/leasedkv_test.go b/etcd/leasedkv_test.go new file mode 100644 index 000000000..288a02052 --- /dev/null +++ b/etcd/leasedkv_test.go @@ -0,0 +1,87 @@ +package etcd + +import ( + "context" + "errors" + "io/ioutil" + "os" + "testing" + "time" + + "github.com/pilosa/pilosa/v2/disco" + "go.etcd.io/etcd/embed" + "go.etcd.io/etcd/etcdserver/api/v3client" +) + +const initVal = "test" +const newVal = "newValue" + +func TestLeasedKv(t *testing.T) { + cfg := embed.NewConfig() + + dir, err := ioutil.TempDir("", "leasedkv-*") + if err != nil { + t.Fatal(err) + } + defer os.RemoveAll(dir) + + cfg.Dir = dir + etcd, err := embed.StartEtcd(cfg) + if err != nil { + t.Fatal(err) + } + + cli := v3client.New(etcd.Server) + + lkv := newLeasedKV(cli, "/test", 1) + + ctx := context.Background() + + err = lkv.Start(initVal) + if err != nil { + t.Fatal(err) + } + + v, err := lkv.Get(ctx) + if err != nil { + t.Fatal(err) + } + + if v != initVal { + t.Fatal("obtained value is not the same as expected. Obtained:", v, "Expected:", initVal) + } + + err = lkv.Set(ctx, "otherValue") + if err != nil { + t.Fatal(err) + } + err = lkv.Set(ctx, "3") + if err != nil { + t.Fatal(err) + } + err = lkv.Set(ctx, newVal) + if err != nil { + t.Fatal(err) + } + + v, err = lkv.Get(ctx) + if err != nil { + t.Fatal(err) + } + + if v != newVal { + t.Fatal("obtained value is not the same as expected. Obtained:", v, "Expected:", initVal) + } + + time.Sleep(5 * time.Second) + + lkv.Stop() + + // we need to wait to force the lease expiration + time.Sleep(5 * time.Second) + + _, err = lkv.Get(ctx) + if err == nil || !errors.Is(err, disco.ErrNoResults) { + t.Fatal("expected error:", disco.ErrNoResults, "obtained:", err) + } +} From 955ab2b5a044fd6ef970529866145feb65a38cc0 Mon Sep 17 00:00:00 2001 From: Antonio Navarro Perez Date: Fri, 5 Mar 2021 13:51:19 +0100 Subject: [PATCH 3/8] Requested changes Signed-off-by: Antonio Navarro Perez --- etcd/embed.go | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index db8e5672f..7e4fd8446 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -78,7 +78,7 @@ type Etcd struct { e *embed.Etcd cli *clientv3.Client - heartBeatLeasedKV, resizeLeasedKV *leasedKV + heartbeatLeasedKV, resizeLeasedKV *leasedKV } func NewEtcd(opt Options, replicas int) *Etcd { @@ -100,8 +100,8 @@ func (e *Etcd) Close() error { e.resizeLeasedKV.Stop() e.resizeLeasedKV = nil } - if e.heartBeatLeasedKV != nil { - e.heartBeatLeasedKV.Stop() + if e.heartbeatLeasedKV != nil { + e.heartbeatLeasedKV.Stop() } e.e.Close() @@ -202,9 +202,9 @@ func (e *Etcd) Start(ctx context.Context) (disco.InitialClusterState, error) { func (e *Etcd) startHeartbeat(ctx context.Context) error { key := heartbeatPrefix + e.e.Server.ID().String() - e.heartBeatLeasedKV = newLeasedKV(e.cli, key, e.options.HeartbeatTTL) + e.heartbeatLeasedKV = newLeasedKV(e.cli, key, e.options.HeartbeatTTL) - if err := e.heartBeatLeasedKV.Start(string(disco.NodeStateStarting)); err != nil { + if err := e.heartbeatLeasedKV.Start(string(disco.NodeStateStarting)); err != nil { return errors.Wrap(err, "startHeartbeat: starting a new heartbeat") } @@ -283,7 +283,7 @@ func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, erro } func (e *Etcd) Started(ctx context.Context) (err error) { - return e.heartBeatLeasedKV.Set(ctx, string(disco.NodeStateStarted)) + return e.heartbeatLeasedKV.Set(ctx, string(disco.NodeStateStarted)) } func (e *Etcd) ID() string { From df0064c55f1cd46b7117339acbd2cfa04908a79a Mon Sep 17 00:00:00 2001 From: Antonio Navarro Perez Date: Fri, 5 Mar 2021 14:03:48 +0100 Subject: [PATCH 4/8] Add godoc. Signed-off-by: Antonio Navarro Perez --- etcd/leasedkv.go | 18 ++++++++++++++---- etcd/leasedkv_test.go | 14 ++++++++++++++ 2 files changed, 28 insertions(+), 4 deletions(-) diff --git a/etcd/leasedkv.go b/etcd/leasedkv.go index 77df1c73c..ad1a46514 100644 --- a/etcd/leasedkv.go +++ b/etcd/leasedkv.go @@ -26,15 +26,19 @@ import ( "go.etcd.io/etcd/clientv3/clientv3util" ) +// leasedKV is an etcd key and value attached to a lease. It can be used to detect if a node went down. +// It will try to renew the lease at any cost after losing it. +// It will recreate the previous existing value for the key again. type leasedKV struct { cli *clientv3.Client cancel context.CancelFunc - mu sync.Mutex - - key, value string - stopped bool + key string ttlSeconds int64 + + mu sync.Mutex + value string // protected by mu + stopped bool // protected by mu } func newLeasedKV(cli *clientv3.Client, key string, ttlSeconds int64) *leasedKV { @@ -45,6 +49,8 @@ func newLeasedKV(cli *clientv3.Client, key string, ttlSeconds int64) *leasedKV { } } +// Start creates the key and attaches it to a lease. +// If the lease cannot be renewed in time, it will try to renew it ad finitum. func (l *leasedKV) Start(initValue string) error { l.mu.Lock() defer l.mu.Unlock() @@ -105,6 +111,8 @@ func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) { } } +// Stop will cancel the lease renewal. +// After calling Stop, this object should be discarded and not used anymore. func (l *leasedKV) Stop() { l.mu.Lock() defer l.mu.Unlock() @@ -116,6 +124,7 @@ func (l *leasedKV) Stop() { l.stopped = true } +// Set will change the specific value for this key. func (l *leasedKV) Set(ctx context.Context, value string) error { l.mu.Lock() defer l.mu.Unlock() @@ -131,6 +140,7 @@ func (l *leasedKV) Set(ctx context.Context, value string) error { return nil } +// Get will obtain the actual value for the key. func (l *leasedKV) Get(ctx context.Context) (string, error) { l.mu.Lock() defer l.mu.Unlock() diff --git a/etcd/leasedkv_test.go b/etcd/leasedkv_test.go index 288a02052..bcd1f06de 100644 --- a/etcd/leasedkv_test.go +++ b/etcd/leasedkv_test.go @@ -1,3 +1,17 @@ +// 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 ( From 85a43c519dbcbdbb19fd26123c66d43ced6c70e4 Mon Sep 17 00:00:00 2001 From: Antonio Navarro Perez Date: Fri, 5 Mar 2021 14:34:53 +0100 Subject: [PATCH 5/8] Fix problem with context errors. Signed-off-by: Antonio Navarro Perez --- etcd/leasedkv.go | 25 +++++++++++++++++-------- 1 file changed, 17 insertions(+), 8 deletions(-) diff --git a/etcd/leasedkv.go b/etcd/leasedkv.go index ad1a46514..b8db8c194 100644 --- a/etcd/leasedkv.go +++ b/etcd/leasedkv.go @@ -89,19 +89,23 @@ func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) { for { kresp, ok := <-ch if !ok { + l.mu.Lock() + if l.stopped { + l.mu.Unlock() return } - l.mu.Lock() - defer l.mu.Unlock() - - retry(1*time.Second, func() error { + if ok := retry(1*time.Second, func() error { return l.create(l.value) - }) + }); !ok { + log.Println("lease cannot be recreated. Key:", l.key) + l.mu.Unlock() + return + } log.Println("lease recreated after a problem. Key:", l.key) - + l.mu.Unlock() return } @@ -162,11 +166,16 @@ func (l *leasedKV) Get(ctx context.Context) (string, error) { return l.value, nil } -func retry(sleep time.Duration, f func() error) { +func retry(sleep time.Duration, f func() error) bool { for { err := f() if err == nil { - return + return true + } + + // sometimes the element in charge of stopping the lease renewal doesn't do it, causing context errors. + if errors.Is(err, context.DeadlineExceeded) { + return false } time.Sleep(sleep) From feea0fa347bde2a08f6bec3680c132041bcd6a0f Mon Sep 17 00:00:00 2001 From: Antonio Navarro Perez Date: Fri, 5 Mar 2021 18:40:00 +0100 Subject: [PATCH 6/8] Requested changes. Signed-off-by: Antonio Navarro Perez --- etcd/leasedkv.go | 32 +++++++++++++++++++------------- 1 file changed, 19 insertions(+), 13 deletions(-) diff --git a/etcd/leasedkv.go b/etcd/leasedkv.go index b8db8c194..82f399270 100644 --- a/etcd/leasedkv.go +++ b/etcd/leasedkv.go @@ -55,34 +55,38 @@ func (l *leasedKV) Start(initValue string) error { l.mu.Lock() defer l.mu.Unlock() - return l.create(initValue) + kaChann, err := l.create(l.value) + if err != nil { + return err + } + + go l.consumeLease(kaChann) + return nil } -func (l *leasedKV) create(initValue string) error { +func (l *leasedKV) create(initValue string) (<-chan *clientv3.LeaseKeepAliveResponse, error) { ctx, cancel := context.WithCancel(context.Background()) l.cancel = cancel leaseResp, err := l.cli.Grant(ctx, l.ttlSeconds) if err != nil { - return errors.Wrap(err, "creating a lease") + return nil, 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) + return nil, errors.Wrapf(err, "creating key %s with value [%s]", l.key, initValue) } kaChann, err := l.cli.KeepAlive(ctx, leaseResp.ID) if err != nil { - return errors.Wrapf(err, "keeping alive the lease for the key %s with value %s", l.key, l.value) + return nil, errors.Wrapf(err, "keeping alive the lease for the key %s with value %s", l.key, l.value) } l.value = initValue - go l.consumeLease(kaChann) - - return nil + return kaChann, nil } func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) { @@ -97,7 +101,13 @@ func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) { } if ok := retry(1*time.Second, func() error { - return l.create(l.value) + kaChann, err := l.create(l.value) + if err != nil { + return err + } + + go l.consumeLease(kaChann) + return nil }); !ok { log.Println("lease cannot be recreated. Key:", l.key) l.mu.Unlock() @@ -108,10 +118,6 @@ func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) { l.mu.Unlock() return } - - if kresp == nil { - continue - } } } From 292794373442ae0ea353decdcde2d925e810f777 Mon Sep 17 00:00:00 2001 From: Antonio Navarro Perez Date: Fri, 5 Mar 2021 18:42:14 +0100 Subject: [PATCH 7/8] Copy paste fix. Signed-off-by: Antonio Navarro Perez --- etcd/leasedkv.go | 2 +- etcd/leasedkv_test.go | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/etcd/leasedkv.go b/etcd/leasedkv.go index 82f399270..3f3bf8a73 100644 --- a/etcd/leasedkv.go +++ b/etcd/leasedkv.go @@ -91,7 +91,7 @@ func (l *leasedKV) create(initValue string) (<-chan *clientv3.LeaseKeepAliveResp func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) { for { - kresp, ok := <-ch + _, ok := <-ch if !ok { l.mu.Lock() diff --git a/etcd/leasedkv_test.go b/etcd/leasedkv_test.go index bcd1f06de..826d05c70 100644 --- a/etcd/leasedkv_test.go +++ b/etcd/leasedkv_test.go @@ -84,7 +84,7 @@ func TestLeasedKv(t *testing.T) { } if v != newVal { - t.Fatal("obtained value is not the same as expected. Obtained:", v, "Expected:", initVal) + t.Fatal("obtained value is not the same as expected. Obtained:", v, "Expected:", newVal) } time.Sleep(5 * time.Second) From 542a98548a36f61c91c9ce8570b71369ab41aaf9 Mon Sep 17 00:00:00 2001 From: Antonio Navarro Perez Date: Fri, 5 Mar 2021 23:34:45 +0100 Subject: [PATCH 8/8] Breaking code at friday afternoon... Signed-off-by: Antonio Navarro Perez --- etcd/leasedkv.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/etcd/leasedkv.go b/etcd/leasedkv.go index 3f3bf8a73..3f0703b64 100644 --- a/etcd/leasedkv.go +++ b/etcd/leasedkv.go @@ -55,7 +55,7 @@ func (l *leasedKV) Start(initValue string) error { l.mu.Lock() defer l.mu.Unlock() - kaChann, err := l.create(l.value) + kaChann, err := l.create(initValue) if err != nil { return err }