refactor the PrimaryNodeID logic

This commit is contained in:
Travis 2021-02-03 20:58:25 -06:00
parent da804ee6d5
commit fbdca3c622
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
3 changed files with 46 additions and 19 deletions

View file

@ -1044,25 +1044,18 @@ func (e *Etcd) RemoveShard(ctx context.Context, index, field string, shard uint6
return nil
}
// Nodes implements the Noder interface.
// Nodes implements the Noder interface. It returns the sorted list of nodes
// based on the etcd 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 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")
}
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
@ -1078,14 +1071,17 @@ func (e *Etcd) nodes(includeMeta bool) []*topology.Node {
// PrimaryNodeID implements the Noder interface.
func (e *Etcd) PrimaryNodeID(hasher topology.Hasher) string {
nodes := e.nodes(false)
return topology.PrimaryNodeID(e.NodeIDs(), hasher)
}
snap := topology.NewClusterSnapshot(topology.NewLocalNoder(nodes), hasher, 1)
primaryNode := snap.PrimaryFieldTranslationNode()
if primaryNode == nil {
return ""
// NodeIDs returns the list of node IDs in the etcd cluster.
func (e *Etcd) NodeIDs() []string {
peers := e.Peers()
ids := make([]string, len(peers))
for i, peer := range peers {
ids[i] = peer.ID
}
return primaryNode.ID
return ids
}
// SetNodes implements the Noder interface as NOP

View file

@ -47,6 +47,25 @@ func NewEmptyLocalNoder() *localNoder {
return &localNoder{}
}
// NewIDNoder is a helper function for wrapping an existing slice of Node IDs
// with something which implements Noder.
func NewIDNoder(ids []string) *localNoder {
nodes := make([]*Node, len(ids))
for i, id := range ids {
node := &Node{
ID: id,
}
nodes[i] = node
}
// Nodes must be sorted.
sort.Sort(ByID(nodes))
return &localNoder{
nodes: nodes,
}
}
// Nodes implements the Noder interface.
func (n *localNoder) Nodes() []*Node {
return n.nodes

View file

@ -280,3 +280,15 @@ func NodePositionByID(nodes []*Node, nodeID string) int {
}
return -1
}
// PrimaryNodeID returns the ID of the primary node, given a list of node IDs
// and a hasher. The order of the node IDs provided does not matter because this
// function will re-order them in a deterministic way.
func PrimaryNodeID(nodeIDs []string, hasher Hasher) string {
snap := NewClusterSnapshot(NewIDNoder(nodeIDs), hasher, 1)
primaryNode := snap.PrimaryFieldTranslationNode()
if primaryNode == nil {
return ""
}
return primaryNode.ID
}