diff --git a/api.go b/api.go index 56435ff6f..2902fde2a 100644 --- a/api.go +++ b/api.go @@ -497,7 +497,10 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string, qcx := api.Txf().NewQcx() defer qcx.Abort() - nodes := api.cluster.shardNodes(indexName, shard) + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + + nodes := snap.ShardNodes(indexName, shard) errCh := make(chan error, len(nodes)) for _, node := range nodes { node := node @@ -619,8 +622,11 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin return errors.Wrap(err, "validating api method") } + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + // Validate that this handler owns the shard. - if !api.cluster.ownsShard(api.Node().ID, indexName, shard) { + if !snap.OwnsShard(api.Node().ID, indexName, shard) { api.server.logger.Printf("node %s does not own shard %d of index %s", api.Node().ID, shard, indexName) return ErrClusterDoesNotOwnShard } @@ -668,7 +674,7 @@ func (api *API) ExportCSV(ctx context.Context, indexName string, fieldName strin } if index.Keys() { - if store := index.TranslateStore(api.cluster.idPartition(indexName, columnID)); store == nil { + if store := index.TranslateStore(snap.IDToShardPartition(indexName, columnID)); store == nil { return errors.Wrap(err, "partition does not exist") } else if colStr, err = store.TranslateID(columnID); err != nil { return errors.Wrap(err, "translating column") @@ -702,7 +708,10 @@ func (api *API) ShardNodes(ctx context.Context, indexName string, shard uint64) return nil, errors.Wrap(err, "validating api method") } - return api.cluster.shardNodes(indexName, shard), nil + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + + return snap.ShardNodes(indexName, shard), nil } // FragmentBlockData is an endpoint for internal usage. It is not guaranteed to @@ -1683,8 +1692,10 @@ func (api *API) LongQueryTime() time.Duration { } func (api *API) validateShardOwnership(indexName string, shard uint64) error { + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) // Validate that this handler owns the shard. - if !api.cluster.ownsShard(api.Node().ID, indexName, shard) { + if !snap.OwnsShard(api.Node().ID, indexName, shard) { api.server.logger.Printf("node %s does not own shard %d of index %s", api.Node().ID, shard, indexName) return ErrClusterDoesNotOwnShard } @@ -2003,7 +2014,10 @@ func (api *API) CreateFieldKeys(ctx context.Context, index, field string, keys . // PrimaryReplicaNodeURL returns the URL of the cluster's primary replica. func (api *API) PrimaryReplicaNodeURL() url.URL { - node := api.cluster.PrimaryReplicaNode() + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(api.cluster.noder, api.cluster.Hasher, api.cluster.ReplicaN) + + node := snap.PrimaryReplicaNode(api.Node().ID) if node == nil { return url.URL{} } diff --git a/cluster.go b/cluster.go index 4544a0034..4a1dcd074 100644 --- a/cluster.go +++ b/cluster.go @@ -666,9 +666,12 @@ func (c *cluster) fragsByHost(idx *Index) fragsByHost { // by creating every combination of field/view specified in `fieldViews` up // for the given set of shards with data. func (c *cluster) fragCombos(idx string, availableShards *roaring.Bitmap, fieldViews viewsByField) fragsByHost { + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + t := make(fragsByHost) _ = availableShards.ForEach(func(i uint64) error { - nodes := c.shardNodes(idx, i) + nodes := snap.ShardNodes(idx, i) for _, n := range nodes { // for each field/view combination: for field, views := range fieldViews { @@ -828,9 +831,13 @@ func (c *cluster) translationNodes(to *cluster) (map[string][]*translationResize m[n.ID] = nil } + // Create a snapshot of the cluster to use for node/partition calculations. + fSnap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + toSnap := topology.NewClusterSnapshot(to.noder, c.Hasher, to.ReplicaN) + for pid := 0; pid < c.partitionN; pid++ { - fNodes := c.partitionNodes(pid) - tNodes := to.partitionNodes(pid) + fNodes := fSnap.PartitionNodes(pid) + tNodes := toSnap.PartitionNodes(pid) // For `to` cluster, we include all nodes containing a // replica for the partition. The source for each replica @@ -888,9 +895,12 @@ func (c *cluster) shardDistributionByIndex(indexName string) map[string]map[stri c.mu.RLock() defer c.mu.RUnlock() + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + for _, shard := range available { - p := c.shardToShardPartition(indexName, shard) - nodes := c.partitionNodes(p) + p := snap.ShardToShardPartition(indexName, shard) + nodes := snap.PartitionNodes(p) dist[nodes[0].ID]["primary-shards"] = append(dist[nodes[0].ID]["primary-shards"], shard) for k := 1; k < len(nodes); k++ { dist[nodes[k].ID]["replica-shards"] = append(dist[nodes[k].ID]["replica-shards"], shard) @@ -931,11 +941,6 @@ func keyToKeyPartition(index, key string, partitionN int) int { return int(h.Sum64() % uint64(partitionN)) } -// idPartition returns the partition that an id belongs to. -func (c *cluster) idPartition(index string, id uint64) int { - return shardToShardPartition(index, id/ShardWidth, c.partitionN) -} - // ShardNodes returns a list of nodes that own a fragment. Safe for concurrent use. func (c *cluster) ShardNodes(index string, shard uint64) []*topology.Node { c.mu.RLock() @@ -960,13 +965,6 @@ func (c *cluster) keyNodes(index, key string) []*topology.Node { return c.partitionNodes(c.Topology.KeyPartition(index, key)) } -// ownsShard returns true if a host owns a fragment. -func (c *cluster) ownsShard(nodeID string, index string, shard uint64) bool { - c.mu.RLock() - defer c.mu.RUnlock() - return topology.Nodes(c.shardNodes(index, shard)).ContainsID(nodeID) -} - // partitionNodes returns a list of nodes that own a partition. unprotected. func (c *cluster) partitionNodes(partitionID int) []*topology.Node { // Default replica count to between one and the number of nodes. @@ -1466,10 +1464,14 @@ func (c *cluster) unprotectedGenerateResizeJobByAction(nodeAction nodeAction) (* j.IDs[node.ID] = true continue } + + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + instr := &ResizeInstruction{ JobID: j.ID, Node: toCluster.unprotectedNodeByID(node.ID), - Coordinator: c.unprotectedCoordinatorNode(), + Coordinator: snap.PrimaryFieldTranslationNode(), Sources: fragmentSourcesByNode[node.ID], TranslationSources: translationSourcesByNode[node.ID], NodeStatus: c.nodeStatus(), // Include the NodeStatus in order to ensure that schema and availableShards are in sync on the receiving node. @@ -1490,7 +1492,10 @@ func (c *cluster) completeCurrentJob(state string) error { } func (c *cluster) unprotectedCompleteCurrentJob(state string) error { - if !c.unprotectedIsCoordinator() { + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + // TODO: this needs to become: IsPrimaryFieldTranslationNode(c.Node.ID) + if !snap.IsCoordinatorNode(c.Node.ID) { return ErrNodeNotCoordinator } if c.currentJob == nil { @@ -2419,16 +2424,19 @@ func (c *cluster) unprotectedPrimaryReplicaNode() *topology.Node { // the case where the local node is not coordinator, then this method will forward the translation // request to the coordinator. func (c *cluster) translateFieldKeys(ctx context.Context, field *Field, keys []string, writable bool) (ids []uint64, err error) { - coordinator := c.coordinatorNode() - if coordinator == nil { + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + + primary := snap.PrimaryFieldTranslationNode() + if primary == nil { return nil, errors.Errorf("translating field(%s/%s) keys(%v) - cannot find coordinator node", field.Index(), field.Name(), keys) } - if c.Node.ID == coordinator.ID { + if c.Node.ID == primary.ID { ids, err = field.TranslateStore().TranslateKeys(keys, writable) } else { // If it's writable, then forward the request to the coordinator. - ids, err = c.InternalClient.TranslateKeysNode(ctx, &coordinator.URI, field.Index(), field.Name(), keys, writable) + ids, err = c.InternalClient.TranslateKeysNode(ctx, &primary.URI, field.Index(), field.Name(), keys, writable) } if err != nil { @@ -2588,15 +2596,18 @@ func (c *cluster) translateFieldIDs(field *Field, ids map[uint64]struct{}) (map[ } func (c *cluster) translateFieldListIDs(field *Field, ids []uint64) (keys []string, err error) { - coordinator := c.coordinatorNode() - if coordinator == nil { + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + + primary := snap.PrimaryFieldTranslationNode() + if primary == nil { return nil, errors.Errorf("translating field(%s/%s) ids(%v) - cannot find coordinator node", field.Index(), field.Name(), ids) } - if c.Node.ID == coordinator.ID { + if c.Node.ID == primary.ID { keys, err = field.TranslateStore().TranslateIDs(ids) } else { - keys, err = c.InternalClient.TranslateIDsNode(context.Background(), &coordinator.URI, field.Index(), field.Name(), ids) + keys, err = c.InternalClient.TranslateIDsNode(context.Background(), &primary.URI, field.Index(), field.Name(), ids) } if err != nil { return nil, errors.Wrapf(err, "translating field(%s/%s) ids(%v)", field.Index(), field.Name(), ids) @@ -2669,10 +2680,13 @@ func (c *cluster) translateIndexKeySet(ctx context.Context, indexName string, ke return nil, ErrIndexNotFound } + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + // Split keys by partition. keysByPartition := make(map[int][]string, c.partitionN) for key := range keySet { - partitionID := c.Topology.KeyPartition(indexName, key) + partitionID := snap.KeyToKeyPartition(indexName, key) keysByPartition[partitionID] = append(keysByPartition[partitionID], key) } @@ -2686,7 +2700,7 @@ func (c *cluster) translateIndexKeySet(ctx context.Context, indexName string, ke g.Go(func() (err error) { var ids []uint64 - primary := c.primaryPartitionNode(partitionID) + primary := snap.PrimaryPartitionNode(partitionID) if primary == nil { return errors.Errorf("translating index(%s) keys(%v) on partition(%d) - cannot find primary node", indexName, keys, partitionID) } @@ -2949,10 +2963,13 @@ func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idS return nil, newNotFoundError(ErrIndexNotFound, indexName) } + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + // Split ids by partition. idsByPartition := make(map[int][]uint64, c.partitionN) for id := range idSet { - partitionID := c.idPartition(indexName, id) + partitionID := snap.IDToShardPartition(indexName, id) idsByPartition[partitionID] = append(idsByPartition[partitionID], id) } @@ -2966,7 +2983,7 @@ func (c *cluster) translateIndexIDSet(ctx context.Context, indexName string, idS g.Go(func() (err error) { var keys []string - primary := c.primaryPartitionNode(partitionID) + primary := snap.PrimaryPartitionNode(partitionID) if primary == nil { return errors.Errorf("translating index(%s) ids(%v) on partition(%d) - cannot find primary node", indexName, ids, partitionID) } diff --git a/executor.go b/executor.go index f38c2a0b5..16706215c 100644 --- a/executor.go +++ b/executor.go @@ -4712,8 +4712,11 @@ func (e *executor) executeClearBitField(ctx context.Context, qcx *Qcx, index str shard := colID / ShardWidth + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN) + ret := false - for _, node := range e.Cluster.shardNodes(index, shard) { + for _, node := range snap.ShardNodes(index, shard) { // Update locally if host matches. if node.ID == e.Node.ID { @@ -5070,7 +5073,10 @@ func (e *executor) executeSetBitField(ctx context.Context, qcx *Qcx, index strin shard := colID / ShardWidth ret := false - for _, node := range e.Cluster.shardNodes(index, shard) { + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN) + + for _, node := range snap.ShardNodes(index, shard) { // Update locally if host matches. if node.ID == e.Node.ID { @@ -5113,7 +5119,10 @@ func (e *executor) executeSetValueField(ctx context.Context, qcx *Qcx, index str shard := colID / ShardWidth ret := false - for _, node := range e.Cluster.shardNodes(index, shard) { + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN) + + for _, node := range snap.ShardNodes(index, shard) { // Update locally if host matches. if node.ID == e.Node.ID { @@ -5157,10 +5166,12 @@ func (e *executor) executeClearValueField(ctx context.Context, qcx *Qcx, index s shard := colID / ShardWidth ret := false - for _, node := range e.Cluster.shardNodes(index, shard) { + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(e.Cluster.noder, e.Cluster.Hasher, e.Cluster.ReplicaN) + + for _, node := range snap.ShardNodes(index, shard) { // Update locally if host matches. if node.ID == e.Node.ID { - idx := e.Holder.Index(index) tx, finisher, err := qcx.GetTx(Txo{Write: writable, Index: idx, Shard: shard}) if err != nil { @@ -5441,9 +5452,15 @@ func (e *executor) remoteExec(ctx context.Context, node *topology.Node, index st func (e *executor) shardsByNode(nodes []*topology.Node, index string, shards []uint64) (map[*topology.Node][]uint64, error) { m := make(map[*topology.Node][]uint64) + // Create a snapshot of the cluster to use for node/partition calculations. + // We use e.Cluster.Nodes() here instead of e.Cluster.noder because we need + // the node states in order to ensure that we don't include an unavailable + // node in the map of nodes to which we distribute the query. + snap := topology.NewClusterSnapshot(topology.NewLocalNoder(e.Cluster.Nodes()), e.Cluster.Hasher, e.Cluster.ReplicaN) + loop: for _, shard := range shards { - for _, node := range e.Cluster.ShardNodes(index, shard) { + for _, node := range snap.ShardNodes(index, shard) { if topology.Nodes(nodes).Contains(node) { m[node] = append(m[node], shard) continue loop diff --git a/holder.go b/holder.go index f3db23fc0..aebbc9deb 100644 --- a/holder.go +++ b/holder.go @@ -1343,6 +1343,10 @@ func (s *holderSyncer) SyncHolder() error { s.mu.Lock() // only allow one instance of SyncHolder to be running at a time defer s.mu.Unlock() ti := time.Now() + + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN) + // Iterate over schema in sorted order. for _, di := range s.Holder.Schema() { // Verify syncer has not closed. @@ -1377,7 +1381,7 @@ func (s *holderSyncer) SyncHolder() error { itr.Seek(0) for shard, eof := itr.Next(); !eof; shard, eof = itr.Next() { // Ignore shards that this host doesn't own. - if !s.Cluster.ownsShard(s.Node.ID, di.Name, shard) { + if !snap.OwnsShard(s.Node.ID, di.Name, shard) { continue } @@ -1539,16 +1543,19 @@ func (s *holderSyncer) resetTranslationSync() error { return errors.Wrap(err, "stop translation sync") } + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN) + // Set read-only flag for all translation stores. - s.setTranslateReadOnlyFlags() + s.setTranslateReadOnlyFlags(snap) // Connect to each node that has a primary for which we are a replica. - if err := s.initializeIndexTranslateReplication(); err != nil { + if err := s.initializeIndexTranslateReplication(snap); err != nil { return errors.Wrap(err, "initialize index translate replication") } // Connect to coordinator to stream field data. - if err := s.initializeFieldTranslateReplication(); err != nil { + if err := s.initializeFieldTranslateReplication(snap); err != nil { return errors.Wrap(err, "initialize field translate replication") } return nil @@ -1619,9 +1626,10 @@ func (s *holderSyncer) stopTranslationSync() error { // setTranslateReadOnlyFlags updates all translation stores to enable or disable // writing new translation keys. Index stores are writable if the node owns the // partition. Field stores are writable if the node is the coordinator. -func (s *holderSyncer) setTranslateReadOnlyFlags() { +func (s *holderSyncer) setTranslateReadOnlyFlags(snap *topology.ClusterSnapshot) { s.Cluster.mu.RLock() - isCoordinator := s.Cluster.unprotectedIsCoordinator() + // TODO: this needs to become: IsPrimaryFieldTranslationNode(s.Cluster.Node.ID) { + isPrimaryFieldTranslator := snap.IsCoordinatorNode(s.Cluster.Node.ID) for _, index := range s.Holder.Indexes() { // There is a race condition here: @@ -1642,8 +1650,8 @@ func (s *holderSyncer) setTranslateReadOnlyFlags() { // // Update: there was another path down to Index.Close(), so // we shrink to lock to be inside index.TranslateStore() now. - for partitionID := 0; partitionID < s.Cluster.partitionN; partitionID++ { - primary := s.Cluster.unprotectedPrimaryPartitionNode(partitionID) + for partitionID := 0; partitionID < snap.PartitionN; partitionID++ { + primary := snap.PrimaryPartitionNode(partitionID) isPrimary := primary != nil && s.Node.ID == primary.ID if ts := index.TranslateStore(partitionID); ts != nil { @@ -1652,7 +1660,7 @@ func (s *holderSyncer) setTranslateReadOnlyFlags() { } for _, field := range index.Fields() { - field.TranslateStore().SetReadOnly(!isCoordinator) + field.TranslateStore().SetReadOnly(!isPrimaryFieldTranslator) } } s.Cluster.mu.RUnlock() @@ -1660,8 +1668,8 @@ func (s *holderSyncer) setTranslateReadOnlyFlags() { // initializeIndexTranslateReplication connects to each node that is the // primary for a partition that we are a replica of. -func (s *holderSyncer) initializeIndexTranslateReplication() error { - for _, node := range s.Cluster.Nodes() { +func (s *holderSyncer) initializeIndexTranslateReplication(snap *topology.ClusterSnapshot) error { + for _, node := range snap.Nodes { // Skip local node. if node.ID == s.Node.ID { continue @@ -1673,8 +1681,8 @@ func (s *holderSyncer) initializeIndexTranslateReplication() error { if !index.Keys() { continue } - for partitionID := 0; partitionID < s.Cluster.partitionN; partitionID++ { - partitionNodes := s.Cluster.partitionNodes(partitionID) + for partitionID := 0; partitionID < snap.PartitionN; partitionID++ { + partitionNodes := snap.PartitionNodes(partitionID) isPrimary := partitionNodes[0].ID == node.ID // remote is primary? isReplica := topology.Nodes(partitionNodes[1:]).ContainsID(s.Node.ID) // local is replica? if !isPrimary || !isReplica { @@ -1713,9 +1721,10 @@ func (s *holderSyncer) initializeIndexTranslateReplication() error { } // initializeFieldTranslateReplication connects the coordinator to stream field data. -func (s *holderSyncer) initializeFieldTranslateReplication() error { +func (s *holderSyncer) initializeFieldTranslateReplication(snap *topology.ClusterSnapshot) error { // Skip if coordinator. - if s.Cluster.isCoordinator() { + // TODO: this needs to become: IsPrimaryFieldTranslationNode(s.Cluster.Node.ID) { + if !snap.IsCoordinatorNode(s.Cluster.Node.ID) { return nil } @@ -1737,9 +1746,9 @@ func (s *holderSyncer) initializeFieldTranslateReplication() error { return nil } - // Connect to coordinator and begin streaming. - coordinator := s.Cluster.coordinatorNode() - rd, err := s.Holder.OpenTranslateReader(context.Background(), coordinator.URI.String(), m) + // Connect to primary and begin streaming. + primary := snap.PrimaryFieldTranslationNode() + rd, err := s.Holder.OpenTranslateReader(context.Background(), primary.URI.String(), m) if err != nil { return err } @@ -1754,6 +1763,9 @@ func (s *holderSyncer) initializeFieldTranslateReplication() error { } func (s *holderSyncer) readIndexTranslateReader(rd TranslateEntryReader) { + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(s.Cluster.noder, s.Cluster.Hasher, s.Cluster.ReplicaN) + for { var entry TranslateEntry if err := rd.ReadEntry(&entry); err != nil { @@ -1769,7 +1781,7 @@ func (s *holderSyncer) readIndexTranslateReader(rd TranslateEntryReader) { } // Apply replication to store. - store := idx.TranslateStore(s.Cluster.Topology.KeyPartition(entry.Index, entry.Key)) + store := idx.TranslateStore(snap.KeyToKeyPartition(entry.Index, entry.Key)) if err := store.ForceSet(entry.ID, entry.Key); err != nil { s.Holder.Logger.Printf("cannot force set index translation data: %d=%q", entry.ID, entry.Key) return @@ -1825,6 +1837,9 @@ func (c *holderCleaner) IsClosing() bool { // CleanHolder compares the holder with the cluster state and removes // any unnecessary fragments and files. func (c *holderCleaner) CleanHolder() error { + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.Cluster.noder, c.Cluster.Hasher, c.Cluster.ReplicaN) + for _, index := range c.Holder.Indexes() { // Verify cleaner has not closed. if c.IsClosing() { @@ -1832,7 +1847,7 @@ func (c *holderCleaner) CleanHolder() error { } // Get the fragments that node is responsible for (based on hash(index, node)). - containedShards := c.Cluster.containsShards(index.Name(), index.AvailableShards(includeRemote), c.Node) + containedShards := snap.ContainsShards(index.Name(), index.AvailableShards(includeRemote), c.Node) // Get the fragments registered in memory. for _, field := range index.Fields() { diff --git a/topology/snapshot.go b/topology/snapshot.go index 2eaa7b0b9..e355ac81a 100644 --- a/topology/snapshot.go +++ b/topology/snapshot.go @@ -138,6 +138,18 @@ func (c *ClusterSnapshot) IsPrimaryFieldTranslationNode(nodeID string) bool { return c.PrimaryFieldTranslationNode().ID == nodeID } +// IsCoordinatorNode returns true if nodeID represents the coordinator +// node responsible for field translation. TODO: this is temporary until +// we transition over to using primary +func (c *ClusterSnapshot) IsCoordinatorNode(nodeID string) bool { + for i := range c.Nodes { + if c.Nodes[i].ID == nodeID && c.Nodes[i].IsCoordinator { + return true + } + } + return false +} + // PrimaryPartitionNode returns the primary node of the given partition. func (c *ClusterSnapshot) PrimaryPartitionNode(partition int) *Node { if nodes := c.PartitionNodes(partition); len(nodes) > 0 {