add PrimaryNodeID() method to Noder interface

This commit is contained in:
Travis 2021-02-02 23:05:08 -06:00
parent 893b5ea6ab
commit f9661b7b81
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
4 changed files with 65 additions and 44 deletions

View file

@ -818,20 +818,6 @@ func (c *cluster) partitionNodes(partitionID int) []*topology.Node {
return nodes
}
func (c *cluster) primaryPartitionNode(partition int) *topology.Node {
c.mu.RLock()
defer c.mu.RUnlock()
return c.unprotectedPrimaryPartitionNode(partition)
}
// unprotectedPrimaryPartition returns tprimary node of partition.
func (c *cluster) unprotectedPrimaryPartitionNode(partition int) *topology.Node {
if nodes := c.partitionNodes(partition); len(nodes) > 0 {
return nodes[0]
}
return nil
}
func (t *Topology) IsPrimary(nodeID string, partitionID int) bool {
primary := t.PrimaryNodeIndex(partitionID)
return nodeID == t.nodeIDs[primary]
@ -1651,6 +1637,11 @@ func (t *Topology) Nodes() []*topology.Node {
return nodes
}
// PrimaryNodeID implements the Noder interface.
func (t *Topology) PrimaryNodeID(topology.Hasher) string {
return ""
}
// SetNodes implements the Noder interface.
func (t *Topology) SetNodes(nodes []*topology.Node) {}
@ -2466,11 +2457,14 @@ func (c *cluster) findIndexKeys(ctx context.Context, indexName string, keys ...s
// TODO: use local replicas to short-circuit network traffic
// Create a snapshot of the cluster to use for node/partition calculations.
snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN)
// Group keys by node.
keysByNode := make(map[*topology.Node][]string)
for partitionID, keys := range keysByPartition {
// Find the primary node for this partition.
primary := c.primaryPartitionNode(partitionID)
primary := snap.PrimaryPartitionNode(partitionID)
if primary == nil {
return nil, errors.Errorf("translating index(%s) keys(%v) on partition(%d) - cannot find primary node", indexName, keys, partitionID)
}
@ -2572,12 +2566,15 @@ func (c *cluster) createIndexKeys(ctx context.Context, indexName string, keys ..
// TODO: use local replicas to short-circuit network traffic
// Create a snapshot of the cluster to use for node/partition calculations.
snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN)
// Group keys by node.
// Delete remote keys from the by-partition map so that it can be used for local translation.
keysByNode := make(map[*topology.Node][]string)
for partitionID, keys := range keysByPartition {
// Find the primary node for this partition.
primary := c.primaryPartitionNode(partitionID)
primary := snap.PrimaryPartitionNode(partitionID)
if primary == nil {
return nil, errors.Errorf("translating index(%s) keys(%v) on partition(%d) - cannot find primary node", indexName, keys, partitionID)
}

View file

@ -1045,18 +1045,24 @@ func (e *Etcd) RemoveShard(ctx context.Context, index, field string, shard uint6
}
// Nodes implements the Noder interface.
func (n *Etcd) Nodes() []*topology.Node {
// If we have looked up nodes within a certain time, then we're going to
// use the cached value for now. This is temporary and will be addressed
// correctly in #1133.
peers := n.Peers()
func (e *Etcd) Nodes() []*topology.Node {
return e.nodes(true)
}
// nodes is a helper function used to get the sorted list of nodes based on the
// etcd peers.
func (e *Etcd) nodes(includeMeta bool) []*topology.Node {
peers := e.Peers()
nodes := make([]*topology.Node, len(peers))
for i, peer := range peers {
node := &topology.Node{}
if meta, err := n.Metadata(context.Background(), peer.ID); err != nil {
log.Println(err, "getting metadata") // TODO: handle this with a logger
} else if err := json.Unmarshal(meta, node); err != nil {
log.Println(err, "unmarshaling json metadata")
if includeMeta {
if meta, err := e.Metadata(context.Background(), peer.ID); err != nil {
log.Println(err, "getting metadata") // TODO: handle this with a logger
} else if err := json.Unmarshal(meta, node); err != nil {
log.Println(err, "unmarshaling json metadata")
}
}
node.ID = peer.ID
@ -1070,16 +1076,28 @@ func (n *Etcd) Nodes() []*topology.Node {
return nodes
}
// PrimaryNodeID implements the Noder interface.
func (e *Etcd) PrimaryNodeID(hasher topology.Hasher) string {
nodes := e.nodes(false)
snap := topology.NewClusterSnapshot(topology.NewLocalNoder(nodes), hasher, 1)
primaryNode := snap.PrimaryFieldTranslationNode()
if primaryNode == nil {
return ""
}
return primaryNode.ID
}
// SetNodes implements the Noder interface as NOP
// (because we can't force to set nodes for etcd).
func (n *Etcd) SetNodes(nodes []*topology.Node) {}
func (e *Etcd) SetNodes(nodes []*topology.Node) {}
// AppendNode implements the Noder interface as NOP
// (because resizer is responsible for adding new nodes).
func (n *Etcd) AppendNode(node *topology.Node) {}
func (e *Etcd) AppendNode(node *topology.Node) {}
// RemoveNode implements the Noder interface as NOP
// (because resizer is responsible for removing existing nodes)
func (n *Etcd) RemoveNode(nodeID string) bool {
func (e *Etcd) RemoveNode(nodeID string) bool {
return false
}

View file

@ -422,7 +422,7 @@ func NewServer(opts ...ServerOption) (*Server, error) {
stator: disco.NopStator,
metadator: disco.NopMetadator,
resizer: disco.NopResizer,
noder: topology.NewLocalNoder(nil),
noder: topology.NewEmptyLocalNoder(),
sharder: disco.NopSharder,
confirmDownRetries: defaultConfirmDownRetries,
@ -565,14 +565,12 @@ func (s *Server) Open() error {
// Set node ID.
s.nodeID = s.disCo.ID()
// TODO we cannot set IsPrimary here because we don't have all the needed info
node := &topology.Node{
ID: s.nodeID,
URI: s.uri,
GRPCURI: s.grpcURI,
State: nodeStateDown,
// TODO set primary
IsPrimary: false,
ID: s.nodeID,
URI: s.uri,
GRPCURI: s.grpcURI,
State: nodeStateDown,
IsPrimary: s.IsPrimary(),
}
// Set metadata for this node.
@ -1036,10 +1034,7 @@ func (s *Server) mergeRemoteStatus(ns *NodeStatus) error {
// IsPrimary returns if this node is primary right now or not.
func (s *Server) IsPrimary() bool {
if primary := s.cluster.PrimaryReplicaNode(); primary != nil {
return s.nodeID == primary.ID
}
return false
return s.nodeID == s.noder.PrimaryNodeID(s.cluster.Hasher)
}
// monitorDiagnostics periodically polls the Pilosa Indexes for cluster info.
@ -1141,7 +1136,7 @@ func (s *Server) monitorRuntime() {
}
func (srv *Server) StartTransaction(ctx context.Context, id string, timeout time.Duration, exclusive bool, remote bool) (*Transaction, error) {
snap := topology.NewClusterSnapshot(srv.cluster, srv.cluster.Hasher, srv.cluster.partitionN)
snap := topology.NewClusterSnapshot(srv.cluster.noder, srv.cluster.Hasher, srv.cluster.partitionN)
node := srv.node()
if !remote && !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 {
return nil, ErrNodeNotCoordinator
@ -1188,7 +1183,7 @@ func (srv *Server) StartTransaction(ctx context.Context, id string, timeout time
}
func (srv *Server) FinishTransaction(ctx context.Context, id string, remote bool) (*Transaction, error) {
snap := topology.NewClusterSnapshot(srv.cluster, srv.cluster.Hasher, srv.cluster.partitionN)
snap := topology.NewClusterSnapshot(srv.cluster.noder, srv.cluster.Hasher, srv.cluster.partitionN)
node := srv.node()
if !remote && !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 {
return nil, ErrNodeNotCoordinator
@ -1218,7 +1213,7 @@ func (srv *Server) FinishTransaction(ctx context.Context, id string, remote bool
}
func (srv *Server) Transactions(ctx context.Context) (map[string]*Transaction, error) {
snap := topology.NewClusterSnapshot(srv.cluster, srv.cluster.Hasher, srv.cluster.partitionN)
snap := topology.NewClusterSnapshot(srv.cluster.noder, srv.cluster.Hasher, srv.cluster.partitionN)
node := srv.node()
if !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 {
return nil, ErrNodeNotCoordinator
@ -1228,7 +1223,7 @@ func (srv *Server) Transactions(ctx context.Context) (map[string]*Transaction, e
}
func (srv *Server) GetTransaction(ctx context.Context, id string, remote bool) (*Transaction, error) {
snap := topology.NewClusterSnapshot(srv.cluster, srv.cluster.Hasher, srv.cluster.partitionN)
snap := topology.NewClusterSnapshot(srv.cluster.noder, srv.cluster.Hasher, srv.cluster.partitionN)
node := srv.node()
if !remote && !snap.IsPrimaryFieldTranslationNode(node.ID) && len(srv.cluster.Nodes()) > 1 {

View file

@ -22,6 +22,7 @@ import (
// nodes in a cluster can be maintained outside of the cluster struct.
type Noder interface {
Nodes() []*Node // Remember: this has to be sorted correctly!!
PrimaryNodeID(hasher Hasher) string
SetNodes([]*Node)
AppendNode(*Node)
RemoveNode(nodeID string) bool
@ -51,6 +52,16 @@ func (n *localNoder) Nodes() []*Node {
return n.nodes
}
// PrimaryNodeID implements the Noder interface.
func (n *localNoder) PrimaryNodeID(hasher Hasher) string {
snap := NewClusterSnapshot(NewLocalNoder(n.nodes), hasher, 1)
primaryNode := snap.PrimaryFieldTranslationNode()
if primaryNode == nil {
return ""
}
return primaryNode.ID
}
// SetNodes implements the Noder interface.
func (n *localNoder) SetNodes(nodes []*Node) {
n.nodes = nodes