renew heartbeat lease if the lease expires while a node is unavailable

This commit is contained in:
Travis 2021-02-25 22:50:53 -06:00
parent d416b88dda
commit ef40bbb617
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
2 changed files with 74 additions and 23 deletions

View file

@ -37,6 +37,7 @@ 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"
"go.etcd.io/etcd/mvcc/mvccpb"
"go.etcd.io/etcd/pkg/types"
@ -217,23 +218,30 @@ func (e *Etcd) Start(ctx context.Context) (disco.InitialClusterState, error) {
}
func (e *Etcd) startHeartbeat() error {
heartbeatID, ctx, heartbeatCancel, err := e.leaseKeepAlive(context.Background(), e.options.HeartbeatTTL)
ctx, heartbeatCancel := context.WithCancel(context.Background())
e.heartbeatCancel = heartbeatCancel
cb := func(heartbeatID clientv3.LeaseID) error {
key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.ClusterStateStarting
if e.e.Config().ClusterState == embed.ClusterStateFlagExisting {
value = disco.ClusterStateResizing
}
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)
}
e.heartbeatID = heartbeatID
return nil
}
_, err := e.leaseKeepAlive(ctx, heartbeatCancel, e.options.HeartbeatTTL, cb)
if err != nil {
return errors.Wrap(err, "startHeartbeat: creates a new hearbeat")
return errors.Wrap(err, "startHeartbeat: creates a new heartbeat")
}
key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.ClusterStateStarting
if e.e.Config().ClusterState == embed.ClusterStateFlagExisting {
value = disco.ClusterStateResizing
}
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)
}
e.heartbeatID, e.heartbeatCancel = heartbeatID, heartbeatCancel
return nil
}
@ -368,7 +376,11 @@ func (e *Etcd) ClusterState(ctx context.Context) (disco.ClusterState, error) {
}
func (e *Etcd) Resize(ctx context.Context) (func([]byte) error, error) {
resizeID, ctx, resizeCancel, err := e.leaseKeepAlive(ctx, e.options.HeartbeatTTL)
ctx, resizeCancel := context.WithCancel(ctx)
cb := func(clientv3.LeaseID) error { return nil }
resizeID, err := e.leaseKeepAlive(ctx, resizeCancel, e.options.HeartbeatTTL, cb)
if err != nil {
return nil, errors.Wrap(err, "Resize: creates a new hearbeat")
}
@ -707,21 +719,24 @@ func (e *Etcd) delKey(ctx context.Context, key string, withPrefix bool) (err 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, ttl int64) (clientv3.LeaseID, context.Context, context.CancelFunc, error) {
ctx, cancelFunc := context.WithCancel(ctx)
func (e *Etcd) leaseKeepAlive(ctx context.Context, cancelFunc context.CancelFunc, ttl int64, cb func(clientv3.LeaseID) error) (clientv3.LeaseID, error) {
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)
return 0, errors.Wrapf(err, "leaseKeepAlive: creates a new lease (TTL: %v)", ttl)
}
keepaliveFunc := func(tick time.Duration) {
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():
@ -733,10 +748,37 @@ func (e *Etcd) leaseKeepAlive(ctx context.Context, ttl int64) (clientv3.LeaseID,
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
return nil
case <-ticker.C:
if _, err = e.cli.KeepAliveOnce(ctx, leaseResp.ID); err != nil {
_, 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 := &hookedClient{Client: v3client.New(e.e.Server)}
var err error
leaseResp, err = cli.Grant(ctx, ttl)
cli.Close()
if err != nil {
cancelFunc()
return errors.Wrapf(err, "leaseKeepAlive: creates a new lease (TTL: %v)", ttl)
}
// 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)
}
}
@ -744,9 +786,17 @@ func (e *Etcd) leaseKeepAlive(ctx context.Context, ttl int64) (clientv3.LeaseID,
}
e.wg.Add(1)
go keepaliveFunc(time.Second)
go func() {
if err := keepaliveFunc(time.Second); err != nil {
log.Printf("leaseKeepAlive: goroutine err: %v\n", err)
}
}()
return leaseResp.ID, ctx, cancelFunc, nil
if err := cb(leaseResp.ID); err != nil {
return 0, errors.Wrap(err, "calling callback")
}
return leaseResp.ID, nil
}
type hookedClient struct {

View file

@ -367,6 +367,7 @@ func NewConfig() *Config {
c.Etcd.Name = ""
c.Etcd.ClusterName = ""
c.Etcd.InitCluster = c.Name + "=" + c.Etcd.LPeerURL
c.Etcd.HeartbeatTTL = 5
return c
}