diff --git a/etcd/cache.go b/etcd/cache.go index dc7b8a539..8005c87ef 100644 --- a/etcd/cache.go +++ b/etcd/cache.go @@ -19,7 +19,6 @@ import ( "sync" "time" - "github.com/pilosa/pilosa/v2/disco" "github.com/pilosa/pilosa/v2/topology" ) @@ -32,27 +31,11 @@ type EtcdWithCache struct { peerMetadataMu sync.RWMutex peerMetadata map[string][]byte - stateMu sync.Mutex // cluster state cache updates peersMu sync.Mutex // peer-list cache updates - nodes []*topology.Node // unmarshalled Node data - nodesTTL int // seconds - nodesLastRequest time.Time // last time requested - nodeStates map[string]nodeState - nodeStateTTL int // seconds - nodeStateFrequency int // max requests per second allowed before using the cache - - clusterStateVal disco.ClusterState - clusterStateTTL int // seconds - clusterStateFrequency int // max requests per second allowed before using the cache - clusterStateLastRequest time.Time - clusterStateLastCache time.Time -} - -type nodeState struct { - val disco.NodeState - lastRequest time.Time - lastCache time.Time + nodes []*topology.Node // unmarshalled Node data + nodesTTL int // seconds + nodesLastRequest time.Time // last time requested } // NewEtcdWithCache returns a new instance of Cache. @@ -60,14 +43,9 @@ func NewEtcdWithCache(opt Options, replicas int) *EtcdWithCache { return &EtcdWithCache{ Etcd: NewEtcd(opt, replicas), - nodeStateTTL: 6, - nodeStateFrequency: 1, - clusterStateTTL: 6, - clusterStateFrequency: 1, - nodesTTL: 6, + nodesTTL: 6, peerMetadata: make(map[string][]byte), - nodeStates: make(map[string]nodeState), } } @@ -88,71 +66,6 @@ func (c *EtcdWithCache) Metadata(ctx context.Context, peerID string) ([]byte, er return v, err } -// ClusterState is a cache wrapper around the Stator.ClusterState method. -func (c *EtcdWithCache) ClusterState(ctx context.Context) (disco.ClusterState, error) { - c.stateMu.Lock() - defer c.stateMu.Unlock() - - now := time.Now() - if now.Sub(c.clusterStateLastCache) > (time.Duration(c.clusterStateTTL)*time.Second) || - now.Sub(c.clusterStateLastRequest) > (time.Second/time.Duration(c.clusterStateFrequency)) { - v, err := c.Etcd.ClusterState(ctx) - if err == nil { - // In order to avoid NodeState() returning a cached value after - // cluster state has changed, we reset the node state caches to - // ensure that the next call to NodeState() returns the latest - // value. And we only need to do this if the cluster state value has - // actually changed. - if c.clusterStateVal != v { - for k, ns := range c.nodeStates { - ns.lastCache = time.Time{} - c.nodeStates[k] = ns - } - } - - c.clusterStateVal = v - c.clusterStateLastCache = now - c.clusterStateLastRequest = now - } - return v, err - } - c.clusterStateLastRequest = now - return c.clusterStateVal, nil -} - -// NodeState is a cache wrapper around the Stator.NodeState method. -func (c *EtcdWithCache) NodeState(ctx context.Context, peerID string) (disco.NodeState, error) { - c.stateMu.Lock() - defer c.stateMu.Unlock() - - ns := c.nodeStates[peerID] - - now := time.Now() - if now.Sub(ns.lastCache) > (time.Duration(c.nodeStateTTL)*time.Second) || - now.Sub(ns.lastRequest) > (time.Second/time.Duration(c.nodeStateFrequency)) { - v, err := c.Etcd.NodeState(ctx, peerID) - if err == nil { - // In order to avoid ClusterState() returning a cached value after a - // node state has changed, we reset the cluster state cache to - // ensure that the next call to ClusterState() returns the latest - // value. And we only need to do this if the node state value has - // actually changed. - if ns.val != v { - c.clusterStateLastCache = time.Time{} - } - - ns.val = v - ns.lastCache = now - ns.lastRequest = now - c.nodeStates[peerID] = ns - } - return v, err - } - ns.lastRequest = now - c.nodeStates[peerID] = ns - return ns.val, nil -} - // Nodes caches the result of the underlying implementation's node list. func (c *EtcdWithCache) Nodes() []*topology.Node { c.peersMu.Lock() diff --git a/etcd/embed.go b/etcd/embed.go index 159e92fcf..64769b8c9 100644 --- a/etcd/embed.go +++ b/etcd/embed.go @@ -36,6 +36,7 @@ import ( "go.etcd.io/etcd/clientv3/clientv3util" "go.etcd.io/etcd/clientv3/concurrency" "go.etcd.io/etcd/embed" + "go.etcd.io/etcd/mvcc" "go.etcd.io/etcd/mvcc/mvccpb" "go.etcd.io/etcd/pkg/types" ) @@ -237,11 +238,11 @@ func (e *Etcd) NodeState(ctx context.Context, peerID string) (disco.NodeState, e } defer cli.Close() - return e.nodeState(ctx, cli, peerID) + return e.nodeState(ctx, peerID) } -func (e *Etcd) nodeState(ctx context.Context, cli *hookedClient, peerID string) (disco.NodeState, error) { - resp, err := cli.Get(ctx, path.Join(resizePrefix, peerID), clientv3.WithCountOnly()) +func (e *Etcd) nodeState(ctx context.Context, peerID string) (disco.NodeState, error) { + resp, err := e.e.Server.KV().Range([]byte(path.Join(resizePrefix, peerID)), nil, mvcc.RangeOptions{Count: true}) if err != nil { return disco.NodeStateUnknown, err } @@ -249,20 +250,20 @@ func (e *Etcd) nodeState(ctx context.Context, cli *hookedClient, peerID string) return disco.NodeStateResizing, nil } - resp, err = cli.Get(ctx, path.Join(heartbeatPrefix, peerID)) + resp, err = e.e.Server.KV().Range([]byte(path.Join(heartbeatPrefix, peerID)), nil, mvcc.RangeOptions{}) if err != nil { return disco.NodeStateUnknown, err } - if len(resp.Kvs) > 1 { + if len(resp.KVs) > 1 { return disco.NodeStateUnknown, disco.ErrTooManyResults } - if len(resp.Kvs) == 0 { + if len(resp.KVs) == 0 { return disco.NodeStateUnknown, disco.ErrNoResults } - return disco.NodeState(resp.Kvs[0].Value), nil + return disco.NodeState(resp.KVs[0].Value), nil } func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, error) { @@ -276,7 +277,7 @@ func (e *Etcd) NodeStates(ctx context.Context) (map[string]disco.NodeState, erro members := e.e.Server.Cluster().Members() for _, member := range members { - s, err := e.nodeState(ctx, cli, member.ID.String()) + s, err := e.nodeState(ctx, member.ID.String()) if err != nil { log.Println("NodeStates get node state", member.ID.String(), err.Error()) } @@ -347,7 +348,7 @@ func (e *Etcd) ClusterState(ctx context.Context) (disco.ClusterState, error) { ) members := e.e.Server.Cluster().Members() for _, m := range members { - ns, err := e.nodeState(ctx, cli, m.ID.String()) + ns, err := e.nodeState(ctx, m.ID.String()) if err != nil { log.Println("ClusterState get node state", err.Error()) continue