From 30d4687a991d635605b3b77e30b9f668b70bb32f Mon Sep 17 00:00:00 2001 From: Travis Date: Fri, 5 Feb 2021 15:16:13 -0600 Subject: [PATCH] remove type Topology --- api.go | 2 +- cluster.go | 474 ++------------------------------------- cluster_internal_test.go | 402 +-------------------------------- server.go | 10 - topology/snapshot.go | 12 +- translate.go | 3 +- utils_internal_test.go | 354 ----------------------------- 7 files changed, 33 insertions(+), 1224 deletions(-) diff --git a/api.go b/api.go index a3d9107a6..b86b11344 100644 --- a/api.go +++ b/api.go @@ -1859,7 +1859,7 @@ func (api *API) Info() serverInfo { StorageBackend: api.holder.txf.TxType(), ReplicaN: api.cluster.ReplicaN, ShardHash: api.cluster.Hasher.Name(), - KeyHash: api.cluster.Topology.Hasher.Name(), + KeyHash: api.cluster.Hasher.Name(), } } diff --git a/cluster.go b/cluster.go index 94a282d7b..039607ec0 100644 --- a/cluster.go +++ b/cluster.go @@ -16,22 +16,14 @@ package pilosa import ( "context" - "encoding/binary" "encoding/json" "fmt" - "hash/fnv" "io" - "io/ioutil" "math/rand" - "os" - "path/filepath" - "sort" "sync" "time" - "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa/v2/disco" - "github.com/pilosa/pilosa/v2/internal" "github.com/pilosa/pilosa/v2/logger" "github.com/pilosa/pilosa/v2/roaring" "github.com/pilosa/pilosa/v2/topology" @@ -104,8 +96,7 @@ type cluster struct { // nolint: maligned maxWritesPerRequest int // Data directory path. - Path string - Topology *Topology + Path string // Distributed Consensus disCo disco.DisCo @@ -751,9 +742,12 @@ func (c *cluster) Nodes() []*topology.Node { } func (c *cluster) AllNodeStates() map[string]string { - c.mu.RLock() - defer c.mu.RUnlock() - return c.Topology.nodeStates + // TODO: is this being used by the UI? + // c.mu.RLock() + // defer c.mu.RUnlock() + // return c.Topology.nodeStates + m := make(map[string]string) + return m } // removeNodeBasicSorted removes a node from the cluster, maintaining the sort @@ -1058,214 +1052,6 @@ func (c *cluster) shardDistributionByIndex(indexName string) map[string]map[stri return dist } -// shardPartition returns the shard-partition that a shard belongs to. -// NOTE: this is DIFFERENT from the key-partition -func (c *cluster) shardToShardPartition(index string, shard uint64) int { - return shardToShardPartition(index, shard, c.partitionN) -} - -func shardToShardPartition(index string, shard uint64, partitionN int) int { - var buf [8]byte - binary.BigEndian.PutUint64(buf[:], shard) - - // Hash the bytes and mod by partition count. - h := fnv.New64a() - _, _ = h.Write([]byte(index)) - _, _ = h.Write(buf[:]) - return int(h.Sum64() % uint64(partitionN)) -} - -// KeyPartition returns the key-partition that a key belongs to. -// NOTE: the key-partition is DIFFERENT from the shard-partition. -func (t *Topology) KeyPartition(index, key string) int { - return keyToKeyPartition(index, key, t.PartitionN) -} - -func keyToKeyPartition(index, key string, partitionN int) int { - // Hash the bytes and mod by partition count. - h := fnv.New64a() - _, _ = h.Write([]byte(index)) - _, _ = h.Write([]byte(key)) - return int(h.Sum64() % uint64(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() - defer c.mu.RUnlock() - return c.shardNodes(index, shard) -} - -// shardNodes returns a list of nodes that own a shard. unprotected -func (c *cluster) shardNodes(index string, shard uint64) []*topology.Node { - return c.partitionNodes(c.shardToShardPartition(index, shard)) -} - -// KeyNodes returns a list of nodes that own a fragment. Safe for concurrent use. -func (c *cluster) KeyNodes(index, key string) []*topology.Node { - c.mu.RLock() - defer c.mu.RUnlock() - return c.keyNodes(index, key) -} - -// keyNodes returns a list of nodes that own a key. unprotected -func (c *cluster) keyNodes(index, key string) []*topology.Node { - return c.partitionNodes(c.Topology.KeyPartition(index, key)) -} - -// 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. - // The replica count can be zero if there are no nodes. - - // Assume that c.nodes may be missing a node that is part of the cluster but not currently present. - // The partition calculation must use the full cluster size in BOTH cases: - // - use len(c.Topology.nodeIDs) instead of len(c.nodes), - // - collect nodes from c.Topology.nodeIDs rather than from c.nodes, - // - when the node is missing, it should be considered, found absent from c.nodes, then omitted from the return slice. - - // Use c.Topology to determine cluster membership when it - // exists and contains data. Otherwise, fall back to using - // c.nodes. The only time c.Topology should be nil is in - // tests. - var useTopology bool - if c.Topology != nil && len(c.Topology.nodeIDs) > 0 { - useTopology = true - } - - cNodes := c.noder.Nodes() - - replicaN := c.ReplicaN - var nodeN int - if useTopology { - nodeN = len(c.Topology.nodeIDs) - } else { - nodeN = len(cNodes) - } - if replicaN > nodeN { - replicaN = nodeN - } else if replicaN == 0 { - replicaN = 1 - } - - // Determine primary owner node. - if c.Topology == nil { - c.Topology = NewTopology(c.Hasher, c.partitionN, c.ReplicaN, c) - } - nodeIndex := c.Topology.PrimaryNodeIndex(partitionID) - if nodeIndex < 0 { - // no nodes anyway - return nil - } - // Collect nodes around the ring. - nodes := make([]*topology.Node, 0, replicaN) - for i := 0; i < replicaN; i++ { - if useTopology { - maybeNodeID := c.Topology.nodeIDs[(nodeIndex+i)%nodeN] - if node := topology.Nodes(cNodes).NodeByID(maybeNodeID); node != nil { - nodes = append(nodes, node) - } - } else { - nodes = append(nodes, cNodes[(nodeIndex+i)%len(cNodes)]) - } - } - - return nodes -} - -func (t *Topology) IsPrimary(nodeID string, partitionID int) bool { - primary := t.PrimaryNodeIndex(partitionID) - return nodeID == t.nodeIDs[primary] -} - -func (t *Topology) PrimaryNodeIndex(partitionID int) (nodeIndex int) { - n := len(t.nodeIDs) - if n == 0 { - if t.cluster != nil { - n = len(t.cluster.noder.Nodes()) - } - } - nodeIndex = t.Hasher.Hash(uint64(partitionID), n) - return -} - -func (t *Topology) GetNonPrimaryReplicas(partitionID int) (nonPrimaryReplicas []string) { - - primary := t.PrimaryNodeIndex(partitionID) - nodeN := len(t.nodeIDs) - - // Collect nodes around the ring. - for i := 1; i < nodeN; i++ { - nodeID := t.nodeIDs[(primary+i)%nodeN] - if i < t.ReplicaN { - nonPrimaryReplicas = append(nonPrimaryReplicas, nodeID) - } - } - return -} - -// the map replicaNodeIDs[nodeID] will have a true value for the primary nodeID, and false for others. -func (t *Topology) GetReplicasForPrimary(primary int) (replicaNodeIDs, nonReplicas map[string]bool) { - if primary < 0 { - // no nodes anyway - return - } - replicaNodeIDs = make(map[string]bool) - nonReplicas = make(map[string]bool) - - nodeN := len(t.nodeIDs) - - // Collect nodes around the ring. - for i := 0; i < nodeN; i++ { - nodeID := t.nodeIDs[(primary+i)%nodeN] - if i < t.ReplicaN { - // mark true if primary - replicaNodeIDs[nodeID] = (i == 0) - } else { - nonReplicas[nodeID] = false - } - } - return -} - -// containsShards is like OwnsShards, but it includes replicas. -func (c *cluster) containsShards(index string, availableShards *roaring.Bitmap, node *topology.Node) []uint64 { - var shards []uint64 - _ = availableShards.ForEach(func(i uint64) error { - p := c.shardToShardPartition(index, i) - // Determine the nodes for partition. - nodes := c.partitionNodes(p) - for _, n := range nodes { - if n.ID == node.ID { - shards = append(shards, i) - } - } - return nil - }) - return shards -} - -func (c *cluster) setup() error { - // Load topology file if it exists. - if err := c.loadTopology(); err != nil { - return errors.Wrap(err, "loading topology") - } - return nil -} - -// open is only used in internal tests. -func (c *cluster) open() error { - err := c.setup() - if err != nil { - return errors.Wrap(err, "setting up cluster") - } - return c.waitForStarted() -} - -func (c *cluster) waitForStarted() error { - return nil -} - func (c *cluster) close() error { // Notify goroutines of closing and wait for completion. close(c.closing) @@ -1478,128 +1264,6 @@ func newResizeJob(existingNodes []*topology.Node, node *topology.Node, action st } } -type nodeIDs []string - -func (n nodeIDs) Len() int { return len(n) } -func (n nodeIDs) Swap(i, j int) { n[i], n[j] = n[j], n[i] } -func (n nodeIDs) Less(i, j int) bool { return n[i] < n[j] } - -// ContainsID returns true if id matches one of the nodesets's IDs. -func (n nodeIDs) ContainsID(id string) bool { - for _, nid := range n { - if nid == id { - return true - } - } - return false -} - -// Topology represents the list of hosts in the cluster. -// Topology now encapsulates all knowledge needed to -// determine the primary node in the replication scheme. -type Topology struct { - mu sync.RWMutex - nodeIDs []string - - clusterID string - - // nodeStates holds the state of each node according to - // the coordinator. Used during startup and data load. - nodeStates map[string]string - - // moved Hasher, PartitionN and ReplicaN - // from cluster for standalone use and comprehension: - - // Hashing algorithm used to assign partitions to nodes. - Hasher topology.Hasher - // The number of partitions in the cluster. - PartitionN int - // The number of replicas a partition has. - ReplicaN int - - // can be nil - cluster *cluster -} - -// NewTopology creates a Topology. -// -// The arguments and members hasher, partitionN, and -// replicaN were refactored out of struct cluster -// to allow pilosa-fsck to load a Topology from -// backup and then compute primaries standalone -- without starting a cluster. -// As pilosa-fsck operates on all backups at once from -// a single cpu, starting a full cluster isn't possible. -// -// The hasher is the Hashing algorithm used to assign partitions to nodes. -// The cluster c should be provided if possible by pilosa code; -// the pilosa-fsck utility won't be able to provide it. -// -// For the cluster size N, the topology gives preference to -// len(t.nodeIDs) before falling back on len(c.nodes). -// -func NewTopology(hasher topology.Hasher, partitionN int, replicaN int, c *cluster) *Topology { - return &Topology{ - Hasher: hasher, - PartitionN: partitionN, - ReplicaN: replicaN, - nodeStates: make(map[string]string), - cluster: c, - } -} - -func (t *Topology) String() string { - return fmt.Sprintf(` -&pilosa.Topology{ - nodeIDs: %v, - clusterID: %v, - nodeStates: %v, - PartitionN: %v, - ReplicaN: %v, -} -`, - t.nodeIDs, - t.clusterID, - t.nodeStates, - t.PartitionN, - t.ReplicaN, - ) -} - -/////////////////////////////////////////// -// Topology implements the Noder interface. - -// Nodes implements the Noder interface. -func (t *Topology) Nodes() []*topology.Node { - nodes := make([]*topology.Node, len(t.nodeIDs)) - for i, nodeID := range t.nodeIDs { - nodes[i] = &topology.Node{ - ID: nodeID, - } - } - 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) {} - -// AppendNode implements the Noder interface. -func (t *Topology) AppendNode(node *topology.Node) {} - -// RemoveNode implements the Noder interface. -func (t *Topology) RemoveNode(nodeID string) bool { - return false -} - -// SetNodeState implements the Noder interface. -func (t *Topology) SetNodeState(nodeID string, state string) {} - -/////////////////////////////////////////// - /////////////////////////////////////////// // Cluster implements the Noder interface. // This is temporary and should be removed once etcd is fully implemented as @@ -1621,66 +1285,6 @@ func (c *cluster) SetNodeState(nodeID string, state string) {} /////////////////////////////////////////// -func (t *Topology) GetNodeIDs() []string { - return t.nodeIDs -} - -// ContainsID returns true if id matches one of the topology's IDs. -func (t *Topology) ContainsID(id string) bool { - t.mu.RLock() - defer t.mu.RUnlock() - return t.containsID(id) -} - -func (t *Topology) containsID(id string) bool { - return nodeIDs(t.nodeIDs).ContainsID(id) -} - -// addID adds the node ID to the topology and returns true if added. -func (t *Topology) addID(nodeID string) bool { - t.mu.Lock() - defer t.mu.Unlock() - if t.containsID(nodeID) { - return false - } - t.nodeIDs = append(t.nodeIDs, nodeID) - - sort.Slice(t.nodeIDs, - func(i, j int) bool { - return t.nodeIDs[i] < t.nodeIDs[j] - }) - - return true -} - -// encode converts t into its internal representation. -func (t *Topology) encode() *internal.Topology { - return encodeTopology(t) -} - -// loadTopology reads the topology for the node. unprotected. -func (c *cluster) loadTopology() error { - buf, err := ioutil.ReadFile(filepath.Join(c.Path, ".topology")) - if os.IsNotExist(err) { - c.Topology = NewTopology(c.Hasher, c.partitionN, c.ReplicaN, c) - return nil - } else if err != nil { - return errors.Wrap(err, "reading file") - } - - var pb internal.Topology - if err := proto.Unmarshal(buf, &pb); err != nil { - return errors.Wrap(err, "unmarshalling") - } - top, err := DecodeTopology(&pb, c.Hasher, c.partitionN, c.ReplicaN, c) - if err != nil { - return errors.Wrap(err, "decoding") - } - c.Topology = top - - return nil -} - func (c *cluster) nodeStatus() *NodeStatus { ns := &NodeStatus{ Node: c.Node, @@ -1974,27 +1578,6 @@ func (c *cluster) translateIndexKeys(ctx context.Context, indexName string, keys return ids, nil } -// The boltdb key translation stores are partitioned, designated by partitionIDs. These -// are shared between replicas, and one node is the primary for -// replication. So with 4 nodes and 3-way replication, each node has 3/4 of -// the translation stores on it. -func (t *Topology) GetPrimaryForColKeyTranslation(index, key string) (primary int) { - partitionID := t.KeyPartition(index, key) - return t.PrimaryNodeIndex(partitionID) -} - -// should match cluster.go:1033 cluster.ownsShard(nodeID, index, shard) -// return Nodes(c.shardNodes(index, shard)).ContainsID(nodeID) -func (t *Topology) GetPrimaryForShardReplication(index string, shard uint64) int { - n := len(t.nodeIDs) - if n == 0 { - return -1 - } - partition := uint64(shardToShardPartition(index, shard, t.PartitionN)) - nodeIndex := t.Hasher.Hash(partition, n) - return nodeIndex -} - func (c *cluster) translateIndexKeySet(ctx context.Context, indexName string, keySet map[string]struct{}, writable bool) (map[string]uint64, error) { keyMap := make(map[string]uint64) @@ -2062,18 +1645,18 @@ func (c *cluster) findIndexKeys(ctx context.Context, indexName string, keys ...s 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 keys { - partitionID := c.Topology.KeyPartition(indexName, key) + partitionID := snap.KeyToKeyPartition(indexName, key) keysByPartition[partitionID] = append(keysByPartition[partitionID], key) } // 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 { @@ -2171,18 +1754,18 @@ func (c *cluster) createIndexKeys(ctx context.Context, indexName string, keys .. return nil, errors.Errorf("can't create index keys on unkeyed index %s", indexName) } + // 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 keys { - partitionID := c.Topology.KeyPartition(indexName, key) + partitionID := snap.KeyToKeyPartition(indexName, key) keysByPartition[partitionID] = append(keysByPartition[partitionID], key) } // 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) @@ -2393,33 +1976,6 @@ type Schema struct { Indexes []*IndexInfo `json:"indexes"` } -func encodeTopology(topology *Topology) *internal.Topology { - if topology == nil { - return nil - } - return &internal.Topology{ - ClusterID: topology.clusterID, - NodeIDs: topology.nodeIDs, - } -} - -// the cluster c is optional but give it if you have it. -func DecodeTopology(topology *internal.Topology, hasher topology.Hasher, partitionN, replicaN int, c *cluster) (*Topology, error) { - if topology == nil { - return nil, nil - } - - t := NewTopology(hasher, partitionN, replicaN, c) - t.clusterID = topology.ClusterID - t.nodeIDs = topology.NodeIDs - sort.Slice(t.nodeIDs, - func(i, j int) bool { - return t.nodeIDs[i] < t.nodeIDs[j] - }) - - return t, nil -} - // CreateShardMessage is an internal message indicating shard creation. type CreateShardMessage struct { Index string diff --git a/cluster_internal_test.go b/cluster_internal_test.go index 68c44a7fb..938c65c0b 100644 --- a/cluster_internal_test.go +++ b/cluster_internal_test.go @@ -15,7 +15,6 @@ package pilosa import ( - "bytes" "fmt" "math/rand" "net" @@ -31,7 +30,6 @@ import ( "github.com/pilosa/pilosa/v2/test/port" "github.com/pilosa/pilosa/v2/testhook" "github.com/pilosa/pilosa/v2/topology" - "github.com/pkg/errors" ) // GlobalPortMap avoids many races and port conflicts when setting @@ -423,13 +421,16 @@ func TestCluster_Owners(t *testing.T) { cNodes := c.noder.Nodes() + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + // Verify nodes are distributed. - if a := c.partitionNodes(0); !reflect.DeepEqual(a, []*topology.Node{cNodes[0], cNodes[1]}) { + if a := snap.PartitionNodes(0); !reflect.DeepEqual(a, []*topology.Node{cNodes[0], cNodes[1]}) { t.Fatalf("unexpected owners: %s", spew.Sdump(a)) } // Verify nodes go around the ring. - if a := c.partitionNodes(2); !reflect.DeepEqual(a, []*topology.Node{cNodes[2], cNodes[0]}) { + if a := snap.PartitionNodes(2); !reflect.DeepEqual(a, []*topology.Node{cNodes[2], cNodes[0]}) { t.Fatalf("unexpected owners: %s", spew.Sdump(a)) } } @@ -440,7 +441,7 @@ func TestCluster_Partition(t *testing.T) { c := newCluster() c.partitionN = partitionN - partitionID := c.shardToShardPartition(index, shard) + partitionID := topology.ShardToShardPartition(index, shard, partitionN) if partitionID < 0 || partitionID >= partitionN { t.Errorf("partition out of range: shard=%d, p=%d, n=%d", shard, partitionID, partitionN) } @@ -483,7 +484,11 @@ func TestCluster_ContainsShards(t *testing.T) { c := NewTestCluster(t, 5) c.ReplicaN = 3 cNodes := c.noder.Nodes() - shards := c.containsShards("test", roaring.NewBitmap(0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10), cNodes[2]) + + // Create a snapshot of the cluster to use for node/partition calculations. + snap := topology.NewClusterSnapshot(c.noder, c.Hasher, c.ReplicaN) + + shards := snap.ContainsShards("test", roaring.NewBitmap(0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10), cNodes[2]) if !reflect.DeepEqual(shards, []uint64{0, 2, 3, 5, 6, 9, 10}) { t.Fatalf("unexpected shars for node's index: %v", shards) @@ -646,368 +651,6 @@ func TestCluster_Coordinator(t *testing.T) { }) } -func TestCluster_Topology(t *testing.T) { - t.Skip("these tests don't really apply anymore; they were meant to tests the cluster and adding topology nodes.") - - c1 := NewTestCluster(t, 1) // automatically creates Node{ID: "node0"} - - const urisCount = 4 - var uris []pnet.URI - if err := port.GetPorts(func(ports []int) error { - for i := 0; i < urisCount; i++ { - uris = append(uris, NewTestURIFromHostPort(fmt.Sprintf("host%d", i), uint16(ports[i]))) - } - return nil - }, urisCount, 10); err != nil { - t.Fatalf("getting ports: %v", err) - } - - node0 := &topology.Node{ID: "node0", URI: uris[0]} - node1 := &topology.Node{ID: "node1", URI: uris[1]} - node2 := &topology.Node{ID: "node2", URI: uris[2]} - nodeinvalid := &topology.Node{ID: "nodeinvalid", URI: uris[3]} - - t.Run("AddNode", func(t *testing.T) { - err := c1.addNode(node1.ID) - if err != nil { - t.Fatal(err) - } - // add the same host. - err = c1.addNode(node1.ID) - if err != nil { - t.Fatal(err) - } - err = c1.addNode(node2.ID) - if err != nil { - t.Fatal(err) - } - - actual := c1.nodeIDs() - expected := []string{node0.ID, node1.ID, node2.ID} - - if !reflect.DeepEqual(actual, expected) { - t.Errorf("expected: %v, but got: %v", expected, actual) - } - }) - - t.Run("ContainsID", func(t *testing.T) { - if !c1.Topology.ContainsID(node1.ID) { - t.Errorf("!ContainsHost error: %v", node1.ID) - } else if c1.Topology.ContainsID(nodeinvalid.ID) { - t.Errorf("ContainsHost error: %v", nodeinvalid.ID) - } - }) -} - -// Ensure that general cluster functionality works as expected. -func TestCluster_ResizeStates(t *testing.T) { - t.Skip("these tests don't really apply anymore; they were meant to tests the cluster startup process using memberlist and a topology file") - t.Run("Single node, no data", func(t *testing.T) { - tc := NewClusterCluster(t, 1) - - // Open TestCluster. - if err := tc.Open(); err != nil { - t.Fatal(err) - } - - node := tc.Clusters[0] - - state, err := node.State() - if err != nil { - t.Fatal(err) - } - - // Ensure that node comes up in state NORMAL. - if state != string(ClusterStateNormal) { - t.Errorf("expected state: %v, but got: %v", ClusterStateNormal, state) - } - - expectedTop := &Topology{ - nodeIDs: []string{node.Node.ID}, - } - - // Verify topology file. - if !reflect.DeepEqual(node.Topology.nodeIDs, expectedTop.nodeIDs) { - t.Errorf("expected topology: %v, but got: %v", expectedTop.nodeIDs, node.Topology.nodeIDs) - } - - // Close TestCluster. - if err := tc.Close(); err != nil { - t.Fatal(err) - } - }) - - t.Run("Single node, in topology", func(t *testing.T) { - tc := NewClusterCluster(t, 0) - if err := tc.addNode(); err != nil { - t.Fatalf("adding node: %v", err) - } - - node := tc.Clusters[0] - - // write topology to data file - top := &Topology{ - nodeIDs: []string{node.Node.ID}, - } - if err := tc.WriteTopology(node.Path, top); err != nil { - t.Fatalf("writing topology: %v", err) - } - - // Open TestCluster. - if err := tc.Open(); err != nil { - t.Fatal(err) - } - - state, err := node.State() - if err != nil { - t.Fatal(err) - } - - // Ensure that node comes up in state NORMAL. - if state != string(ClusterStateNormal) { - t.Errorf("expected state: %v, but got: %v", ClusterStateNormal, state) - } - - // Close TestCluster. - if err := tc.Close(); err != nil { - t.Fatal(err) - } - }) - - t.Run("Single node, not in topology", func(t *testing.T) { - tc := NewClusterCluster(t, 0) - if err := tc.addNode(); err != nil { - t.Fatalf("adding node: %v", err) - } - - node := tc.Clusters[0] - - // write topology to data file - top := &Topology{ - nodeIDs: []string{"some-other-host"}, - } - if err := tc.WriteTopology(node.Path, top); err != nil { - t.Fatalf("writing topology: %v", err) - } - - // Open TestCluster. - expected := "coordinator node0 is not in topology: [some-other-host]" - err := tc.Open() - if err == nil || errors.Cause(err).Error() != expected { - t.Errorf("did not receive expected error, got: %s", errors.Cause(err).Error()) - } - - // Close TestCluster. - if err := tc.Close(); err != nil { - t.Fatal(err) - } - }) - - t.Run("Multiple nodes, no data", func(t *testing.T) { - tc := NewClusterCluster(t, 0) - if err := tc.addNode(); err != nil { - t.Fatalf("adding node: %v", err) - } - - // Open TestCluster. - if err := tc.Open(); err != nil { - t.Fatalf("opening cluster: %v", err) - } - - if err := tc.addNode(); err != nil { - t.Fatalf("adding node: %v", err) - } - - node0 := tc.Clusters[0] - state0, err := node0.State() - if err != nil { - t.Fatal(err) - } - - node1 := tc.Clusters[1] - state1, err := node1.State() - if err != nil { - t.Fatal(err) - } - - // Ensure that nodes comes up in state NORMAL. - if state0 != string(ClusterStateNormal) { - t.Errorf("expected node0 state: %v, but got: %v", ClusterStateNormal, state0) - } else if state1 != string(ClusterStateNormal) { - t.Errorf("expected node1 state: %v, but got: %v", ClusterStateNormal, state1) - } - - expectedTop := &Topology{ - nodeIDs: []string{node0.Node.ID, node1.Node.ID}, - } - - // Verify topology file. - if !reflect.DeepEqual(node0.Topology.nodeIDs, expectedTop.nodeIDs) { - t.Errorf("expected node0 topology: %v, but got: %v", expectedTop.nodeIDs, node0.Topology.nodeIDs) - } else if !reflect.DeepEqual(node1.Topology.nodeIDs, expectedTop.nodeIDs) { - t.Errorf("expected node1 topology: %v, but got: %v", expectedTop.nodeIDs, node1.Topology.nodeIDs) - } - - // Close TestCluster. - if err := tc.Close(); err != nil { - t.Fatal(err) - } - }) - - t.Run("Multiple nodes, in/not in topology", func(t *testing.T) { - tc := NewClusterCluster(t, 0) - if err := tc.addNode(); err != nil { - t.Fatalf("adding node: %v", err) - } - node0 := tc.Clusters[0] - - // write topology to data file - top := &Topology{ - nodeIDs: []string{"node0", "node2"}, - } - if err := tc.WriteTopology(node0.Path, top); err != nil { - t.Fatalf("writing topology: %v", err) - } - - // Open TestCluster. - if err := tc.Open(); err != nil { - t.Fatalf("opening cluster: %v", err) - } - - state0, err := node0.State() - if err != nil { - t.Fatal(err) - } - - // Ensure that node is in state STARTING before the other node joins. - if state0 != string(ClusterStateStarting) { - t.Errorf("expected node0 state: %v, but got: %v", ClusterStateStarting, state0) - } - - if err := tc.addNode(); err != nil { - t.Fatalf("adding node: %v", err) - } - node1 := tc.Clusters[1] - state1, err := node1.State() - if err != nil { - t.Fatal(err) - } - - // Ensure that node comes up in state NORMAL. - if state0 != string(ClusterStateNormal) { - t.Errorf("expected node0 state: %v, but got: %v", ClusterStateNormal, state0) - } else if state1 != string(ClusterStateNormal) { - t.Errorf("expected node2 state: %v, but got: %v", ClusterStateNormal, state1) - } - - // Close TestCluster. - if err := tc.Close(); err != nil { - t.Fatal(err) - } - }) - - t.Run("Multiple nodes, with data", func(t *testing.T) { - tc := NewClusterCluster(t, 0) - if err := tc.addNode(); err != nil { - t.Fatalf("adding node: %v", err) - } - node0 := tc.Clusters[0] - - // Open TestCluster. - if err := tc.Open(); err != nil { - t.Fatal(err) - } - - // Close TestCluster with defer. - defer func() { - if err := tc.Close(); err != nil { - t.Fatal(err) - } - }() - - // Add Bit Data to node0. - if err := tc.CreateField("i", "f", OptFieldTypeDefault()); err != nil { - t.Fatalf("creating field: %v", err) - } - // Each tc.SetBit starts and commits its own Tx. - if err := tc.SetBit("i", "f", 1, 101, nil); err != nil { - t.Fatalf("setting bit: %v", err) - } - if err := tc.SetBit("i", "f", 1, ShardWidth+1, nil); err != nil { - t.Fatalf("setting bit: %v", err) - } - - // Before starting the resize, get the CheckSum to use for - // comparison later. - node0Field := node0.holder.Field("i", "f") - node0View := node0Field.view("standard") - node0Fragment := node0View.Fragment(1) - node0Checksum, err := node0Fragment.Checksum() - if err != nil { - t.Fatal(err) - } - - idx0 := node0.holder.Index("i") - if idx0 == nil { - t.Fatal(`idx0 was nil, could not retrieve Index("i")`) - } - - // addNode needs to block until the resize process has completed. - if err := tc.addNode(); err != nil { - t.Fatalf("adding node: %v", err) - } - - node1 := tc.Clusters[1] - - state1, err := node1.State() - if err != nil { - t.Fatal(err) - } - - state0, err := node0.State() - if err != nil { - t.Fatal(err) - } - - // Ensure that nodes come up in state NORMAL. - if state0 != string(ClusterStateNormal) { - t.Errorf("expected node0 state: %v, but got: %v", ClusterStateNormal, state0) - } else if state1 != string(ClusterStateNormal) { - t.Errorf("expected node1 state: %v, but got: %v", ClusterStateNormal, state1) - } - // INVAR: after node1.State() is normal, the rebalancing should have been done. - - expectedTop := &Topology{ - nodeIDs: []string{node0.Node.ID, node1.Node.ID}, - } - - // Verify topology file. - if !reflect.DeepEqual(node0.Topology.nodeIDs, expectedTop.nodeIDs) { - t.Errorf("expected node0 topology: %v, but got: %v", expectedTop.nodeIDs, node0.Topology.nodeIDs) - } else if !reflect.DeepEqual(node1.Topology.nodeIDs, expectedTop.nodeIDs) { - t.Errorf("expected node1 topology: %v, but got: %v", expectedTop.nodeIDs, node1.Topology.nodeIDs) - } - - // Bits - // Verify that node-1 contains the fragment (i/f/standard/1) transferred from node-0. - node1Field := node1.holder.Field("i", "f") - node1View := node1Field.view("standard") - node1Fragment := node1View.Fragment(1) - - idx1 := node1.holder.Index("i") - if idx1 == nil { - t.Fatal(`idx1 was nil, could not retrieve Index("i")`) - } - - // Ensure checksums are the same. - if chksum, err := node1Fragment.Checksum(); err != nil { - t.Fatal(err) - } else if !bytes.Equal(chksum, node0Checksum) { - t.Fatalf("expected standard view checksum to match: %x - %x", chksum, node0Checksum) - } - }) -} - func TestAE(t *testing.T) { t.Run("AbortDoesn'tBlockUninitialized", func(t *testing.T) { c := newCluster() @@ -1067,26 +710,3 @@ func TestAE(t *testing.T) { } }) } - -func TestCluster_GetNonPrimaryReplicas(t *testing.T) { - c := newCluster() - c.ReplicaN = 3 - topo := NewTopology(c.Hasher, c.partitionN, c.ReplicaN, c) - c.Topology = topo - nNodes := 4 - for i := 0; i < nNodes; i++ { - nodeID := fmt.Sprintf("node%d", i) - c.noder.AppendNode(&topology.Node{ - ID: nodeID, - URI: NewTestURI("http", fmt.Sprintf("host%d", i), uint16(0)), - }) - c.Topology.addID(nodeID) - } - - partitionID := 256 - nonPrimes := topo.GetNonPrimaryReplicas(partitionID) - m := len(nonPrimes) - if m != c.ReplicaN-1 { - t.Fatalf("expected 2 non primes, got %v", m) - } -} diff --git a/server.go b/server.go index 74ebb76c2..b43e61de6 100644 --- a/server.go +++ b/server.go @@ -585,16 +585,6 @@ func (s *Server) Open() error { s.syncer.Closing = s.closing s.syncer.Stats = s.holder.Stats.WithTags("component:HolderSyncer") - err = s.cluster.setup() - if err != nil { - return errors.Wrap(err, "setting up cluster") - } - - // Open Cluster management. - if err := s.cluster.waitForStarted(); err != nil { - return errors.Wrap(err, "opening Cluster") - } - // Open holder. if err := s.holder.Open(); err != nil { return errors.Wrap(err, "opening Holder") diff --git a/topology/snapshot.go b/topology/snapshot.go index fc1a2d83f..b7d865335 100644 --- a/topology/snapshot.go +++ b/topology/snapshot.go @@ -71,13 +71,11 @@ func NewClusterSnapshot(noder Noder, hasher Hasher, replicas int) *ClusterSnapsh // ShardToShardPartition returns the shard-partition that the given shard // belongs to. NOTE: This is DIFFERENT from the key-partition. func (c *ClusterSnapshot) ShardToShardPartition(index string, shard uint64) int { - return dedupShardToShardPartition(index, shard, c.PartitionN) + return ShardToShardPartition(index, shard, c.PartitionN) } -// dedupShardToShardParition would ideally be called `shardToShardPartition`, but since -// we can't put this into it's own package yet (see the TODO below about import loops), -// that name conflicts with a function that already exists in the `pilosa` package. -func dedupShardToShardPartition(index string, shard uint64, partitionN int) int { +// ShardToShardParition ... +func ShardToShardPartition(index string, shard uint64, partitionN int) int { var buf [8]byte binary.BigEndian.PutUint64(buf[:], shard) @@ -238,14 +236,12 @@ func (c *ClusterSnapshot) PrimaryForColKeyTranslation(index, key string) (primar } // TODO: update this comment -// should match cluster.go:1033 cluster.ownsShard(nodeID, index, shard) -// return Nodes(c.shardNodes(index, shard)).ContainsID(nodeID) func (c *ClusterSnapshot) PrimaryForShardReplication(index string, shard uint64) int { n := len(c.Nodes) if n == 0 { return -1 } - partition := uint64(dedupShardToShardPartition(index, shard, c.PartitionN)) + partition := uint64(ShardToShardPartition(index, shard, c.PartitionN)) nodeIndex := c.Hasher.Hash(partition, n) return nodeIndex } diff --git a/translate.go b/translate.go index 71bf45b4d..dced91507 100644 --- a/translate.go +++ b/translate.go @@ -22,6 +22,7 @@ import ( "io/ioutil" "sync" + "github.com/pilosa/pilosa/v2/topology" "github.com/pkg/errors" ) @@ -178,7 +179,7 @@ func GenerateNextPartitionedID(index string, prev uint64, partitionID, partition // Try to use the next ID if it is in the same partition. // Otherwise find ID in next shard that has a matching partition. for id := prev + 1; ; id += ShardWidth { - if shardToShardPartition(index, id/ShardWidth, partitionN) == partitionID { + if topology.ShardToShardPartition(index, id/ShardWidth, partitionN) == partitionID { return id } } diff --git a/utils_internal_test.go b/utils_internal_test.go index 26414ad11..f0f4ef7d4 100644 --- a/utils_internal_test.go +++ b/utils_internal_test.go @@ -17,18 +17,13 @@ package pilosa import ( "bytes" "fmt" - "io/ioutil" - "path/filepath" - "sync" "testing" "time" - "github.com/gogo/protobuf/proto" pnet "github.com/pilosa/pilosa/v2/net" "github.com/pilosa/pilosa/v2/roaring" "github.com/pilosa/pilosa/v2/testhook" "github.com/pilosa/pilosa/v2/topology" - "github.com/pkg/errors" ) // utilities used by tests @@ -72,7 +67,6 @@ func NewTestCluster(tb testing.TB, n int) *cluster { c.ReplicaN = 1 c.Hasher = NewTestModHasher() c.Path = path - c.Topology = NewTopology(c.Hasher, c.partitionN, c.ReplicaN, c) for i := 0; i < n; i++ { c.noder.AppendNode(&topology.Node{ @@ -113,352 +107,6 @@ func (*TestModHasher) Hash(key uint64, n int) int { return int(key) % n } func (*TestModHasher) Name() string { return "mod" } -// ClusterCluster represents a cluster of test nodes, each of which -// has a Cluster. -// ClusterCluster implements Broadcaster interface. -type ClusterCluster struct { - Clusters []*cluster - - common *commonClusterSettings - - mu sync.RWMutex - resizing bool - resizeDone chan struct{} - tb testing.TB -} - -type commonClusterSettings struct { - Nodes []*topology.Node -} - -func (t *ClusterCluster) CreateIndex(name string) error { - for _, c := range t.Clusters { - if _, err := c.holder.CreateIndexIfNotExists(name, IndexOptions{}); err != nil { - return err - } - } - return nil -} - -func (t *ClusterCluster) CreateIndexWithOpt(name string, opt IndexOptions) error { - for _, c := range t.Clusters { - if _, err := c.holder.CreateIndexIfNotExists(name, opt); err != nil { - return err - } - } - return nil -} - -func (t *ClusterCluster) CreateField(index, field string, opts FieldOption) error { - for _, c := range t.Clusters { - idx, err := c.holder.CreateIndexIfNotExists(index, IndexOptions{}) - if err != nil { - return err - } - if _, err := idx.CreateField(field, opts); err != nil { - return err - } - } - return nil -} - -func (t *ClusterCluster) SetBit(index, field string, rowID, colID uint64, x *time.Time) error { - // Determine which node should receive the SetBit. - c0 := t.Clusters[0] // use the first node's cluster to determine shard location. - shard := colID / ShardWidth - nodes := c0.shardNodes(index, shard) - - for _, node := range nodes { - c := t.clusterByID(node.ID) - if c == nil { - continue - } - f := c.holder.Field(index, field) - if f == nil { - return fmt.Errorf("index/field does not exist: %s/%s", index, field) - } - - if err := func() error { - idx := c.holder.Index(f.index) - shard := colID / ShardWidth - tx := idx.holder.txf.NewTx(Txo{Write: writable, Index: idx, Shard: shard}) - if tx != nil { - defer tx.Rollback() - } - - if _, err := f.SetBit(tx, rowID, colID, x); err != nil { - return err - } else if err := tx.Commit(); err != nil { - return err - } - return nil - }(); err != nil { - return err - } - } - - return nil -} - -func (t *ClusterCluster) clusterByID(id string) *cluster { - for _, c := range t.Clusters { - if c.Node.ID == id { - return c - } - } - return nil -} - -// addNode adds a node to the cluster and (potentially) starts a resize job. -func (t *ClusterCluster) addNode() error { - return nil -} - -// WriteTopology writes the given topology to disk. -func (t *ClusterCluster) WriteTopology(path string, top *Topology) error { - if buf, err := proto.Marshal(top.encode()); err != nil { - return err - } else if err := ioutil.WriteFile(filepath.Join(path, ".topology"), buf, 0666); err != nil { - return err - } - return nil -} - -func (t *ClusterCluster) addCluster(i int, saveTopology bool) (*cluster, error) { - id := fmt.Sprintf("node%d", i) - uri := NewTestURI("http", fmt.Sprintf("host%d", i), uint16(0)) - - node := &topology.Node{ - ID: id, - URI: uri, - } - - // add URI to common - //t.common.NodeIDs = append(t.common.NodeIDs, id) - //sort.Sort(t.common.NodeIDs) - - // add node to common - t.common.Nodes = append(t.common.Nodes, node) - - // create node-specific temp directory - path, err := testhook.TempDirInDir(t.tb, *TempDir, fmt.Sprintf("pilosa-cluster-node-%d-", i)) - if err != nil { - return nil, err - } - - // holder - h := NewHolder(path, nil) - - // cluster - c := newCluster() - c.ReplicaN = 1 - c.Hasher = NewTestModHasher() - c.Path = path - c.partitionN = topology.DefaultPartitionN - c.Topology = NewTopology(c.Hasher, c.partitionN, c.ReplicaN, c) - c.holder = h - c.Node = node - // c.Coordinator = t.common.Nodes[0].ID // the first node is the coordinator - c.broadcaster = t.broadcaster(c) - - // add nodes - if saveTopology { - for _, n := range t.common.Nodes { - if err := c.addNode(n.ID); err != nil { - return nil, err - } - } - } - - // Add this node to the ClusterCluster. - t.Clusters = append(t.Clusters, c) - - return c, nil -} - -// NewClusterCluster returns a new instance of test.Cluster. -func NewClusterCluster(tb testing.TB, n int) *ClusterCluster { - - tc := &ClusterCluster{ - common: &commonClusterSettings{}, - tb: tb, - } - - // add clusters - for i := 0; i < n; i++ { - _, err := tc.addCluster(i, true) - if err != nil { - panic(err) - } - } - return tc -} - -// Open opens all clusters in the test cluster. -func (t *ClusterCluster) Open() error { - for _, c := range t.Clusters { - if err := c.open(); err != nil { - return err - } - if err := c.holder.Open(); err != nil { - return err - } - } - return nil -} - -// Close closes all clusters in the test cluster. -func (t *ClusterCluster) Close() error { - for _, c := range t.Clusters { - err := c.close() - if err != nil { - return err - } - // Make sure open indexes get shut down too. we wouldn't do - // this normally for a cluster, but we want to for test cases. - c.holder.Close() - } - return nil -} - -type bcast struct { - t *ClusterCluster - c *cluster -} - -func (b bcast) SendSync(m Message) error { - switch obj := m.(type) { - case *ClusterStatus: - b.t.mu.RLock() - if obj.State == string(ClusterStateNormal) && b.t.resizing { - close(b.t.resizeDone) - } - b.t.mu.RUnlock() - } - return nil -} - -func (t *ClusterCluster) broadcaster(c *cluster) broadcaster { - return bcast{ - t: t, - c: c, - } -} - -// SendAsync is a test implemenetation of Broadcaster SendAsync method. -func (bcast) SendAsync(Message) error { - return nil -} - -// SendTo is a test implementation of Broadcaster SendTo method. -func (b bcast) SendTo(to *topology.Node, m Message) error { - switch obj := m.(type) { - case *ResizeInstruction: - err := b.t.FollowResizeInstruction(obj) - if err != nil { - return err - } - case *ClusterStatus: - b.t.mu.RLock() - if obj.State == string(ClusterStateNormal) && b.t.resizing { - close(b.t.resizeDone) - } - b.t.mu.RUnlock() - default: - panic(fmt.Sprintf("message not handled:\n%#v\n", obj)) - } - return nil -} - -// FollowResizeInstruction is a version of cluster.followResizeInstruction used for testing. -func (t *ClusterCluster) FollowResizeInstruction(instr *ResizeInstruction) error { - // Prepare the return message. - complete := &ResizeInstructionComplete{ - JobID: instr.JobID, - Node: instr.Node, - Error: "", - } - - // Stop processing on any error. - if err := func() error { - - // figure out which node it was meant for, then call the operation on that cluster - // basically need to mimic this: client.RetrieveShardFromURI(context.Background(), src.Index, src.Field, src.View, src.Shard, srcURI) - instrNode := instr.Node - destCluster := t.clusterByID(instrNode.ID) - - // Sync the schema received in the resize instruction. - if err := destCluster.holder.applySchema(instr.NodeStatus.Schema); err != nil { - return err - } - - // Sync available shards. - for k, is := range instr.NodeStatus.Indexes { - _ = k - for _, fs := range is.Fields { - f := destCluster.holder.Field(is.Name, fs.Name) - - // if we don't know about a field locally, log an error because - // fields should be created and synced prior to shard creation - if f == nil { - continue - } - if err := f.AddRemoteAvailableShards(fs.AvailableShards); err != nil { - return errors.Wrap(err, "adding remote available shards") - } - } - } - - for _, src := range instr.Sources { - srcCluster := t.clusterByID(src.Node.ID) - - srcFragment := srcCluster.holder.fragment(src.Index, src.Field, src.View, src.Shard) - destFragment := destCluster.holder.fragment(src.Index, src.Field, src.View, src.Shard) - if destFragment == nil { - // Create fragment on destination if it doesn't exist. - f := destCluster.holder.Field(src.Index, src.Field) - v := f.view(src.View) - var err error - destFragment, err = v.CreateFragmentIfNotExists(src.Shard) - if err != nil { - return err - } - } - - // this is the *test* version of a network call, transferring fragments between - // nodes in a cluster. So it is allowed to be kind of a hack. - - // there will be two -rbfdb directories/databases, we need to copy - // from src to dest the fragment. This simulates sending the fragment over the network. - srcIdx := srcCluster.holder.Index(src.Index) - srctx := srcIdx.holder.txf.NewTx(Txo{Write: !writable, Index: srcIdx, Fragment: srcFragment, Shard: srcFragment.shard}) - - destIdx := destCluster.holder.Index(src.Index) - - desttx := destIdx.holder.txf.NewTx(Txo{Write: writable, Index: destIdx, Fragment: destFragment, Shard: destFragment.shard}) - - citer, _, err := srctx.ContainerIterator(src.Index, src.Field, src.View, src.Shard, 0) - panicOn(err) - d := destFragment - for citer.Next() { - ckey, c := citer.Value() - err := desttx.PutContainer(d.index(), d.field(), d.view(), d.shard, ckey, c) - panicOn(err) - } - citer.Close() - panicOn(desttx.Commit()) - srctx.Rollback() - } - - return nil - }(); err != nil { - complete.Error = err.Error() - } - - node := instr.Primary - return bcast{t: t}.SendTo(node, complete) -} - var _ = NewTestClusterWithReplication // happy linter func NewTestClusterWithReplication(tb testing.TB, nNodes, nReplicas, partitionN int) (c *cluster, cleaner func()) { @@ -478,7 +126,6 @@ func NewTestClusterWithReplication(tb testing.TB, nNodes, nReplicas, partitionN c.Hasher = &topology.Jmphasher{} c.Path = path c.partitionN = partitionN - c.Topology = NewTopology(c.Hasher, c.partitionN, c.ReplicaN, c) for i := 0; i < nNodes; i++ { nodeID := fmt.Sprintf("node%d", i) @@ -486,7 +133,6 @@ func NewTestClusterWithReplication(tb testing.TB, nNodes, nReplicas, partitionN ID: nodeID, URI: NewTestURI("http", fmt.Sprintf("host%d", i), uint16(0)), }) - c.Topology.addID(nodeID) } cNodes := c.noder.Nodes()