Merge pull request #1513 from ajnavarro/disco/improve-leased-keys

[CORE-312] Improve leased keys.
This commit is contained in:
Matthew Jaffee 2021-03-05 16:47:54 -06:00 • committed by GitHub
commit 863e57d5d1
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
3 changed files with 316 additions and 167 deletions

View file

@ -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 {

191
etcd/leasedkv.go Normal file
View file

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

101
etcd/leasedkv_test.go Normal file
View file

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