diff --git a/etcd/embed.go b/etcd/embed.go index ee0fd1e55..7e4fd8446 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 := heartbeatPrefix + e.e.Server.ID().String() + 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(string(disco.NodeStateStarting)); 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(""); 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..3f0703b64 --- /dev/null +++ b/etcd/leasedkv.go @@ -0,0 +1,191 @@ +// 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" + "log" + "sync" + "time" + + "github.com/pilosa/pilosa/v2/disco" + "github.com/pkg/errors" + "go.etcd.io/etcd/clientv3" + "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 + + 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 { + return &leasedKV{ + cli: cli, + key: key, + ttlSeconds: ttlSeconds, + } +} + +// 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() + + kaChann, err := l.create(initValue) + if err != nil { + return err + } + + go l.consumeLease(kaChann) + return nil +} + +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 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 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 nil, errors.Wrapf(err, "keeping alive the lease for the key %s with value %s", l.key, l.value) + } + + l.value = initValue + + return kaChann, nil +} + +func (l *leasedKV) consumeLease(ch <-chan *clientv3.LeaseKeepAliveResponse) { + for { + _, ok := <-ch + if !ok { + l.mu.Lock() + + if l.stopped { + l.mu.Unlock() + return + } + + if ok := retry(1*time.Second, func() error { + 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() + return + } + + log.Println("lease recreated after a problem. Key:", l.key) + l.mu.Unlock() + return + } + } +} + +// 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() + + if l.cancel != nil { + l.cancel() + } + + 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() + + 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 +} + +// 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() + + 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 +} + +func retry(sleep time.Duration, f func() error) bool { + for { + err := f() + if err == nil { + 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) + + log.Println("retrying after error:", err) + } +} diff --git a/etcd/leasedkv_test.go b/etcd/leasedkv_test.go new file mode 100644 index 000000000..826d05c70 --- /dev/null +++ b/etcd/leasedkv_test.go @@ -0,0 +1,101 @@ +// 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" + "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:", newVal) + } + + 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) + } +}