remove type Topology

This commit is contained in:
Travis 2021-02-05 15:16:13 -06:00
parent 002aee63e3
commit 30d4687a99
No known key found for this signature in database
GPG key ID: 37080CC2042BA34E
7 changed files with 33 additions and 1224 deletions

2
api.go
View file

@ -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(),
}
}

View file

@ -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

View file

@ -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)
}
}

View file

@ -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")

View file

@ -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
}

View file

@ -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
}
}

View file

@ -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()