From 6f7d748c8a6134b22c7456d01e10e1710905e033 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Wed, 3 Mar 2021 13:37:20 +0100 Subject: [PATCH 1/4] Always Put states in Txn --- etcd/embed.go | 157 ++++++++++++++++++++------------------------------ 1 file changed, 62 insertions(+), 95 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index 5ba20d04e..b0e742913 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -34,10 +34,8 @@ import ( "go.etcd.io/etcd/clientv3/clientv3util" "go.etcd.io/etcd/clientv3/concurrency" "go.etcd.io/etcd/embed" - "go.etcd.io/etcd/etcdserver/api/membership" "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" ) @@ -102,6 +100,10 @@ func NewEtcd(opt Options, replicas int) *Etcd { replicas: replicas, wg: &sync.WaitGroup{}, } + + if e.options.HeartbeatTTL == 0 { + e.options.HeartbeatTTL = 5 // seconds + } return e } @@ -217,14 +219,19 @@ func (e *Etcd) startHeartbeat() error { e.heartbeatCancel = heartbeatCancel cb := func(heartbeatID clientv3.LeaseID) error { - key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.ClusterStateStarting + key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.NodeStateStarting if e.e.Config().ClusterState == embed.ClusterStateFlagExisting { - value = disco.ClusterStateResizing + value = disco.NodeStateResizing + } else if e.lm.started { + value = disco.NodeStateStarted } - if _, err := e.cli.Put(ctx, key, string(value), clientv3.WithLease(heartbeatID)); err != nil { + if _, err := e.cli.Txn(ctx). + Then(clientv3.OpPut(key, string(value), clientv3.WithLease(heartbeatID))). + Commit(); err != nil { + heartbeatCancel() - return errors.Wrapf(err, "startHeartbeat: puts a key-value (%s, %s) with lease (%v)", key, value, heartbeatID) + return errors.Wrapf(err, "startHeartbeat: txn puts a key-value (%s, %s) with lease (%v)", key, value, heartbeatID) } e.heartbeatID = heartbeatID @@ -232,7 +239,7 @@ func (e *Etcd) startHeartbeat() error { return nil } - _, err := e.leaseKeepAlive(ctx, heartbeatCancel, e.options.HeartbeatTTL, cb) + _, err := e.leaseKeepAlive(ctx, heartbeatCancel, cb) if err != nil { return errors.Wrap(err, "startHeartbeat: creates a new heartbeat") } @@ -241,43 +248,26 @@ func (e *Etcd) startHeartbeat() error { } func (e *Etcd) NodeState(ctx context.Context, peerID string) (disco.NodeState, error) { - if state, err := e.nodeStateFast(ctx, peerID); err == nil && state == disco.NodeStateStarted { - return disco.NodeStateStarted, nil - } - - states, err := e.NodeStates(ctx) - if err != nil { - return "", err - } - - state, ok := states[peerID] - if !ok { - return disco.NodeStateUnknown, nil - } - - return state, nil + return e.nodeState(ctx, peerID) } -func (e *Etcd) nodeStateFast(ctx context.Context, peerID string) (disco.NodeState, error) { - kv := e.e.Server.KV() - resp, err := kv.Range([]byte(path.Join(resizePrefix, peerID)), nil, mvcc.RangeOptions{Count: true}) +func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, error) { + resp, err := e.cli.Txn(ctx). + If(clientv3util.KeyMissing(path.Join(resizePrefix, peerID))). + Then(clientv3.OpGet(path.Join(heartbeatPrefix, peerID))). + Commit() if err != nil { return disco.NodeStateUnknown, err } - if resp.Count > 0 { + + if !resp.Succeeded { return disco.NodeStateResizing, nil } - resp, err = kv.Range([]byte(path.Join(heartbeatPrefix, peerID)), nil, mvcc.RangeOptions{}) - if err != nil { - return disco.NodeStateUnknown, err - } - kvs := resp.KVs - + kvs := resp.Responses[0].GetResponseRange().Kvs if len(kvs) > 1 { return disco.NodeStateUnknown, disco.ErrTooManyResults } - if len(kvs) == 0 { return disco.NodeStateUnknown, disco.ErrNoResults } @@ -287,10 +277,6 @@ func (e *Etcd) nodeStateFast(ctx context.Context, peerID string) (disco.NodeStat func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, error) { members := e.e.Server.Cluster().Members() - if states := e.nodeStatesFast(ctx, members); states != nil { - return states, nil - } - ops := make([]clientv3.Op, 2*(len(members))) for i, member := range members { peerID := member.ID.String() @@ -298,34 +284,28 @@ func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, erro ops[2*i+1] = clientv3.OpGet(path.Join(heartbeatPrefix, peerID)) } -doTxn: resp, err := e.cli.Txn(ctx).Then(ops...).Commit() if err != nil { return nil, err } - if !resp.Succeeded { - goto doTxn - } out := make(map[string]disco.NodeState, len(members)) for i, member := range members { peerID := member.ID.String() - switch resp.Responses[2*i].GetResponseRange().Count { - case 0: - case 1: + if resp.Responses[2*i].GetResponseRange().Count > 0 { // This node is processing a resize operation. out[peerID] = disco.NodeStateResizing continue - default: - return nil, disco.ErrTooManyResults } - switch resp := resp.Responses[2*i+1].GetResponseRange(); len(resp.Kvs) { + + kvs := resp.Responses[2*i+1].GetResponseRange().Kvs + switch len(kvs) { case 0: // The node has not reported a state. out[peerID] = disco.NodeStateUnknown case 1: // The node has reported its state. - out[peerID] = disco.NodeState(resp.Kvs[0].Value) + out[peerID] = disco.NodeState(kvs[0].Value) default: return nil, disco.ErrTooManyResults } @@ -334,24 +314,13 @@ doTxn: return out, nil } -func (e *Etcd) nodeStatesFast(ctx context.Context, members []*membership.Member) map[string]disco.NodeState { - out := make(map[string]disco.NodeState, len(members)) - for _, member := range members { - peerID := member.ID.String() - state, err := e.nodeStateFast(ctx, peerID) - if err != nil || state != disco.NodeStateStarted { - return nil - } - - out[peerID] = disco.NodeStateStarted - } - - return out -} - func (e *Etcd) Started(ctx context.Context) (err error) { key, value := heartbeatPrefix+e.e.Server.ID().String(), disco.NodeStateStarted - if _, err = e.cli.Put(ctx, key, string(value), clientv3.WithLease(e.heartbeatID)); err == nil { + _, 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 @@ -381,8 +350,13 @@ func (e *Etcd) IsLeader() bool { func (e *Etcd) Leader() *disco.Peer { id := e.e.Server.Leader() - m := e.e.Server.Cluster().Member(id) - return &disco.Peer{ID: id.String(), URL: m.PickPeerURL()} + peer := &disco.Peer{ID: id.String()} + + if m := e.e.Server.Cluster().Member(id); m != nil { + peer.URL = m.PickPeerURL() + } + + return peer } func (e *Etcd) ClusterState(ctx context.Context) (disco.ClusterState, error) { @@ -437,7 +411,7 @@ func (e *Etcd) Resize(ctx context.Context) (func([]byte) error, error) { cb := func(clientv3.LeaseID) error { return nil } - resizeID, err := e.leaseKeepAlive(ctx, resizeCancel, e.options.HeartbeatTTL, cb) + resizeID, err := e.leaseKeepAlive(ctx, resizeCancel, cb) if err != nil { return nil, errors.Wrap(err, "Resize: creates a new hearbeat") } @@ -707,7 +681,9 @@ func (e *Etcd) DeleteView(ctx context.Context, indexName, fieldName, name string } func (e *Etcd) putKey(ctx context.Context, key, val string, opts ...clientv3.OpOption) error { - if _, err := e.cli.Put(ctx, key, val, opts...); err != nil { + if _, err := e.cli.Txn(ctx). + Then(clientv3.OpPut(key, val, opts...)). + Commit(); err != nil { return errors.Wrapf(err, "putKey: Put(%s, %s)", key, val) } @@ -736,11 +712,14 @@ func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { } func (e *Etcd) getKeyWithPrefix(ctx context.Context, key string) ([]string, [][]byte, error) { - resp, err := e.cli.Get(ctx, key, clientv3.WithPrefix()) + + resp, err := e.cli.Txn(ctx). + Then(clientv3.OpGet(key, clientv3.WithPrefix())). + Commit() if err != nil { return nil, nil, err } - kvs := resp.Kvs + kvs := resp.Responses[0].GetResponseRange().Kvs var ( keys []string @@ -756,30 +735,18 @@ func (e *Etcd) getKeyWithPrefix(ctx context.Context, key string) ([]string, [][] } func (e *Etcd) keyExists(ctx context.Context, key string) (bool, error) { - if ok, err := e.keyExistsFast(ctx, key); err == nil && ok { - return true, nil - } - - resp, err := e.cli.Txn(ctx).Then(clientv3.OpGet(key, clientv3.WithCountOnly())).Commit() + resp, err := e.cli.Txn(ctx). + If(clientv3util.KeyExists(key)). + Then(clientv3.OpGet(key, clientv3.WithCountOnly())). + Commit() if err != nil { return false, err } - if resp.Responses[0].GetResponseRange().Count > 0 { - return true, nil + if !resp.Succeeded { + return false, nil } - return false, nil -} -func (e *Etcd) keyExistsFast(ctx context.Context, key string) (bool, error) { - kv := e.e.Server.KV() - resp, err := kv.Range([]byte(key), nil, mvcc.RangeOptions{Count: true}) - if err != nil { - return false, err - } - if resp.Count > 0 { - return true, nil - } - return false, nil + return resp.Responses[0].GetResponseRange().Count > 0, nil } func (e *Etcd) delKey(ctx context.Context, key string, withPrefix bool) (err error) { @@ -794,11 +761,11 @@ 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, cancelFunc context.CancelFunc, ttl int64, cb func(clientv3.LeaseID) error) (clientv3.LeaseID, error) { - leaseResp, err := e.cli.Grant(ctx, ttl) +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: %v)", ttl) + return 0, errors.Wrapf(err, "leaseKeepAlive: creates a new lease (TTL: %d s.)", e.options.HeartbeatTTL) } keepaliveFunc := func(tick time.Duration) error { @@ -818,7 +785,7 @@ func (e *Etcd) leaseKeepAlive(ctx context.Context, cancelFunc context.CancelFunc // 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(ttl)) + 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 { @@ -835,11 +802,11 @@ func (e *Etcd) leaseKeepAlive(ctx context.Context, cancelFunc context.CancelFunc // TODO: should this close/reset e.cli instead? cli := v3client.New(e.e.Server) var err error - leaseResp, err = cli.Grant(ctx, ttl) + leaseResp, err = cli.Grant(ctx, e.options.HeartbeatTTL) cli.Close() if err != nil { cancelFunc() - return errors.Wrapf(err, "leaseKeepAlive: creates a new lease (TTL: %v)", ttl) + return errors.Wrapf(err, "leaseKeepAlive: creates a new lease (TTL: %d s.)", e.options.HeartbeatTTL) } // Call the callback. From a14baf8c15be026262877cf8075b2b8be4fa4165 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Wed, 3 Mar 2021 14:22:44 +0100 Subject: [PATCH 2/4] waitForStatus for cluster test --- http/client.go | 30 +++++++++++++++++++++++++++ internal/clustertests/cluster_test.go | 29 +++++++++++++++++++++++--- 2 files changed, 56 insertions(+), 3 deletions(-) diff --git a/http/client.go b/http/client.go index ce9c60214..53482ddeb 100644 --- a/http/client.go +++ b/http/client.go @@ -104,6 +104,36 @@ func (c *InternalClient) maxShardByIndex(ctx context.Context) (map[string]uint64 return rsp.Standard, nil } +func (c *InternalClient) Status(ctx context.Context) (string, error) { + span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.Status") + defer span.Finish() + + // Execute request against the host. + u := c.defaultURI.Path("/status") + + // Build request. + req, err := http.NewRequest("GET", u, nil) + if err != nil { + return "", errors.Wrap(err, "creating request") + } + + req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) + req.Header.Set("Accept", "application/json") + + // Execute request. + resp, err := c.executeRequest(req.WithContext(ctx)) + if err != nil { + return "", err + } + defer resp.Body.Close() + + var rsp getStatusResponse + if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil { + return "", fmt.Errorf("json decode: %s", err) + } + return rsp.State, nil +} + // SchemaNode returns all index and field schema information from the specified // node. func (c *InternalClient) SchemaNode(ctx context.Context, uri *pnet.URI, views bool) ([]*pilosa.IndexInfo, error) { diff --git a/internal/clustertests/cluster_test.go b/internal/clustertests/cluster_test.go index c156e9125..f4edf0298 100644 --- a/internal/clustertests/cluster_test.go +++ b/internal/clustertests/cluster_test.go @@ -89,9 +89,8 @@ func TestClusterStuff(t *testing.T) { t.Fatalf("waiting on pumba pause cmd: %v", err) } - // TODO change the sleep to wait for status to return to NORMAL - need support in internal client for getting status t.Log("done with pause, waiting for stability") - time.Sleep(time.Second * 20) + waitForStatus(t, cli1, "NORMAL", 30, time.Second) t.Log("done waiting for stability") // Check query results from each node. @@ -105,5 +104,29 @@ func TestClusterStuff(t *testing.T) { } } }) - +} + +func waitForStatus(t *testing.T, c *picli.InternalClient, status string, n int, sleep time.Duration) { + t.Helper() + + for i := 0; i < n; i++ { + s, err := c.Status(context.TODO()) + if err != nil { + t.Logf("Status (%d/%d): %v (sleep: %s)\\n", i, n, err, sleep.String()) + } else { + t.Logf("Status (%d/%d): %s (sleep: %s)\n", i, n, s, sleep.String()) + } + if s == status { + return + } + time.Sleep(sleep) + } + + s, err := c.Status(context.TODO()) + if err != nil { + t.Fatalf("querying status: %v", err) + } + if status != s { + t.Fatalf("waited %d %v for status: %v, got: %v", n, sleep, status, s) + } } From fa293ba6c351b6a2966ec63c7350dd135178daa1 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Wed, 3 Mar 2021 16:11:24 +0100 Subject: [PATCH 3/4] Address PR comments --- etcd/embed.go | 44 ++++++++++++++++----------- internal/clustertests/cluster_test.go | 10 +++--- 2 files changed, 33 insertions(+), 21 deletions(-) diff --git a/etcd/embed.go b/etcd/embed.go index b0e742913..ee0fd1e55 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -264,13 +264,17 @@ func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, e return disco.NodeStateResizing, nil } - kvs := resp.Responses[0].GetResponseRange().Kvs - if len(kvs) > 1 { - return disco.NodeStateUnknown, disco.ErrTooManyResults + if len(resp.Responses) == 0 { + return disco.NodeStateUnknown, disco.ErrNoResults } + + kvs := resp.Responses[0].GetResponseRange().Kvs if len(kvs) == 0 { return disco.NodeStateUnknown, disco.ErrNoResults } + if len(kvs) > 1 { + return disco.NodeStateUnknown, disco.ErrTooManyResults + } return disco.NodeState(kvs[0].Value), nil } @@ -698,12 +702,11 @@ func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { return nil, err } - kvs := resp.Responses[0].GetResponseRange().Kvs - - if !resp.Succeeded { - return nil, errors.New("tx failed") + if len(resp.Responses) == 0 { + return nil, errors.New("key does not exist") } + kvs := resp.Responses[0].GetResponseRange().Kvs if len(kvs) == 0 { return nil, errors.New("key does not exist") } @@ -711,24 +714,28 @@ func (e *Etcd) getKeyBytes(ctx context.Context, key string) ([]byte, error) { return kvs[0].Value, nil } -func (e *Etcd) getKeyWithPrefix(ctx context.Context, key string) ([]string, [][]byte, error) { - +func (e *Etcd) getKeyWithPrefix(ctx context.Context, key string) (keys []string, values [][]byte, err error) { resp, err := e.cli.Txn(ctx). Then(clientv3.OpGet(key, clientv3.WithPrefix())). Commit() if err != nil { return nil, nil, err } + + if len(resp.Responses) == 0 { + return nil, nil, errors.New("key does not exist") + } + kvs := resp.Responses[0].GetResponseRange().Kvs + if len(kvs) == 0 { + return nil, nil, nil + } - var ( - keys []string - values [][]byte - ) - - for _, kv := range kvs { - keys = append(keys, string(kv.Key)) - values = append(values, kv.Value) + keys = make([]string, len(kvs)) + values = make([][]byte, len(kvs)) + for i, kv := range kvs { + keys[i] = string(kv.Key) + values[i] = kv.Value } return keys, values, nil @@ -746,6 +753,9 @@ func (e *Etcd) keyExists(ctx context.Context, key string) (bool, error) { return false, nil } + if len(resp.Responses) == 0 { + return false, nil + } return resp.Responses[0].GetResponseRange().Count > 0, nil } diff --git a/internal/clustertests/cluster_test.go b/internal/clustertests/cluster_test.go index f4edf0298..2f443642c 100644 --- a/internal/clustertests/cluster_test.go +++ b/internal/clustertests/cluster_test.go @@ -22,6 +22,7 @@ import ( "time" "github.com/pilosa/pilosa/v2" + "github.com/pilosa/pilosa/v2/disco" picli "github.com/pilosa/pilosa/v2/http" ) @@ -90,7 +91,7 @@ func TestClusterStuff(t *testing.T) { } t.Log("done with pause, waiting for stability") - waitForStatus(t, cli1, "NORMAL", 30, time.Second) + waitForStatus(t, cli1, string(disco.ClusterStateNormal), 30, time.Second) t.Log("done waiting for stability") // Check query results from each node. @@ -112,9 +113,9 @@ func waitForStatus(t *testing.T, c *picli.InternalClient, status string, n int, for i := 0; i < n; i++ { s, err := c.Status(context.TODO()) if err != nil { - t.Logf("Status (%d/%d): %v (sleep: %s)\\n", i, n, err, sleep.String()) + t.Logf("Status (try %d/%d): %v (retrying in %s)", i, n, err, sleep.String()) } else { - t.Logf("Status (%d/%d): %s (sleep: %s)\n", i, n, s, sleep.String()) + t.Logf("Status (try %d/%d): %s (retrying in %s)", i, n, s, sleep.String()) } if s == status { return @@ -127,6 +128,7 @@ func waitForStatus(t *testing.T, c *picli.InternalClient, status string, n int, t.Fatalf("querying status: %v", err) } if status != s { - t.Fatalf("waited %d %v for status: %v, got: %v", n, sleep, status, s) + waited := time.Duration(n) * sleep + t.Fatalf("waited %s for status: %s, got: %s", waited.String(), status, s) } } From d311b0cac400e83e8fc24b2f3de2afd28718cb7a Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Kuba=20Podg=C3=B3rski?= Date: Wed, 3 Mar 2021 23:47:47 +0100 Subject: [PATCH 4/4] Comment Status function + make waitForStatus more generic --- http/client.go | 63 ++++++++++++++------------- internal/clustertests/cluster_test.go | 8 ++-- 2 files changed, 37 insertions(+), 34 deletions(-) diff --git a/http/client.go b/http/client.go index 53482ddeb..3adce376a 100644 --- a/http/client.go +++ b/http/client.go @@ -104,36 +104,6 @@ func (c *InternalClient) maxShardByIndex(ctx context.Context) (map[string]uint64 return rsp.Standard, nil } -func (c *InternalClient) Status(ctx context.Context) (string, error) { - span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.Status") - defer span.Finish() - - // Execute request against the host. - u := c.defaultURI.Path("/status") - - // Build request. - req, err := http.NewRequest("GET", u, nil) - if err != nil { - return "", errors.Wrap(err, "creating request") - } - - req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) - req.Header.Set("Accept", "application/json") - - // Execute request. - resp, err := c.executeRequest(req.WithContext(ctx)) - if err != nil { - return "", err - } - defer resp.Body.Close() - - var rsp getStatusResponse - if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil { - return "", fmt.Errorf("json decode: %s", err) - } - return rsp.State, nil -} - // SchemaNode returns all index and field schema information from the specified // node. func (c *InternalClient) SchemaNode(ctx context.Context, uri *pnet.URI, views bool) ([]*pilosa.IndexInfo, error) { @@ -2108,3 +2078,36 @@ func (c *InternalClient) ImportFieldKeys(ctx context.Context, uri *pnet.URI, ind defer resp.Body.Close() return nil } + +// Status function is just a public function for this particular implementation of InternalClient. +// It's not require by pilosa.InternalClient interface. +// The function returns pilosa cluster state as a string ("NORMAL", "DEGRADED", "DOWN", "RESIZING", ...) +func (c *InternalClient) Status(ctx context.Context) (string, error) { + span, ctx := tracing.StartSpanFromContext(ctx, "InternalClient.Status") + defer span.Finish() + + // Execute request against the host. + u := c.defaultURI.Path("/status") + + // Build request. + req, err := http.NewRequest("GET", u, nil) + if err != nil { + return "", errors.Wrap(err, "creating request") + } + + req.Header.Set("User-Agent", "pilosa/"+pilosa.Version) + req.Header.Set("Accept", "application/json") + + // Execute request. + resp, err := c.executeRequest(req.WithContext(ctx)) + if err != nil { + return "", err + } + defer resp.Body.Close() + + var rsp getStatusResponse + if err := json.NewDecoder(resp.Body).Decode(&rsp); err != nil { + return "", fmt.Errorf("json decode: %s", err) + } + return rsp.State, nil +} diff --git a/internal/clustertests/cluster_test.go b/internal/clustertests/cluster_test.go index 2f443642c..94e0f7b10 100644 --- a/internal/clustertests/cluster_test.go +++ b/internal/clustertests/cluster_test.go @@ -91,7 +91,7 @@ func TestClusterStuff(t *testing.T) { } t.Log("done with pause, waiting for stability") - waitForStatus(t, cli1, string(disco.ClusterStateNormal), 30, time.Second) + waitForStatus(t, cli1.Status, string(disco.ClusterStateNormal), 30, time.Second) t.Log("done waiting for stability") // Check query results from each node. @@ -107,11 +107,11 @@ func TestClusterStuff(t *testing.T) { }) } -func waitForStatus(t *testing.T, c *picli.InternalClient, status string, n int, sleep time.Duration) { +func waitForStatus(t *testing.T, stator func(context.Context) (string, error), status string, n int, sleep time.Duration) { t.Helper() for i := 0; i < n; i++ { - s, err := c.Status(context.TODO()) + s, err := stator(context.TODO()) if err != nil { t.Logf("Status (try %d/%d): %v (retrying in %s)", i, n, err, sleep.String()) } else { @@ -123,7 +123,7 @@ func waitForStatus(t *testing.T, c *picli.InternalClient, status string, n int, time.Sleep(sleep) } - s, err := c.Status(context.TODO()) + s, err := stator(context.TODO()) if err != nil { t.Fatalf("querying status: %v", err) }