mirror of
https://github.com/featurebasedb/featurebase.git
synced 2026-10-07 11:27:50 +00:00
Merge pull request #929 from travisturner/cluster-resize-tests
Add Cluster resize tests.
This commit is contained in:
commit
638e328455
9 changed files with 631 additions and 60 deletions
|
|
@ -87,7 +87,7 @@ func TestClient_MultiNode(t *testing.T) {
|
|||
}
|
||||
}
|
||||
if !ownsNum {
|
||||
t.Fatalf("Trying to use slice %d on host %s, but it doesn't own that slice. It owns %s", num, s[i].Host(), owns)
|
||||
t.Fatalf("Trying to use slice %d on host %s, but it doesn't own that slice. It owns %v", num, s[i].Host(), owns)
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
57
cluster.go
57
cluster.go
|
|
@ -727,8 +727,8 @@ func (c *Cluster) CompleteCurrentJob(state string) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
// followResizeInstruction is run by any node that receives a ResizeInstruction.
|
||||
func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) {
|
||||
// FollowResizeInstruction is run by any node that receives a ResizeInstruction.
|
||||
func (c *Cluster) FollowResizeInstruction(instr *internal.ResizeInstruction) error {
|
||||
go func() {
|
||||
// Prepare the return message.
|
||||
complete := &internal.ResizeInstructionComplete{
|
||||
|
|
@ -740,31 +740,9 @@ func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) {
|
|||
// Stop processing on any error.
|
||||
if err := func() error {
|
||||
|
||||
// TODO: move this schema creation code to a method on Holder.
|
||||
// Sync the schema received in the resize instruction.
|
||||
// Create indexes that don't exist.
|
||||
for _, index := range instr.Schema.Indexes {
|
||||
opt := IndexOptions{}
|
||||
idx, err := c.Holder.CreateIndexIfNotExists(index.Name, opt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// Create frames that don't exist.
|
||||
for _, f := range index.Frames {
|
||||
opt := decodeFrameOptions(f.Meta)
|
||||
frame, err := idx.CreateFrameIfNotExists(f.Name, *opt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// Create views that don't exist.
|
||||
for _, v := range f.Views {
|
||||
_, err := frame.CreateViewIfNotExists(v)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
// TODO: Create inputDefinitions that don't exist.
|
||||
if err := c.Holder.ApplySchema(instr.Schema); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Create a client for calling remote nodes.
|
||||
|
|
@ -775,7 +753,7 @@ func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) {
|
|||
|
||||
// Request each source file in ResizeSources.
|
||||
for _, src := range instr.Sources {
|
||||
fmt.Printf("\n**** Get slice %d for index %s from host %s ****\n\n", src.Slice, src.Index, src.URI)
|
||||
c.logger().Printf("\n**** Get slice %d for index %s from host %s ****\n\n", src.Slice, src.Index, src.URI)
|
||||
|
||||
srcURI := decodeURI(src.URI)
|
||||
|
||||
|
|
@ -828,6 +806,7 @@ func (c *Cluster) followResizeInstruction(instr *internal.ResizeInstruction) {
|
|||
c.logger().Printf("sending resizeInstructionComplete error: err=%s", err)
|
||||
}
|
||||
}()
|
||||
return nil
|
||||
}
|
||||
|
||||
func (c *Cluster) MarkResizeInstructionComplete(complete *internal.ResizeInstructionComplete) error {
|
||||
|
|
@ -921,11 +900,8 @@ func (j *ResizeJob) setState(state string) {
|
|||
|
||||
// Run distributes ResizeInstructions.
|
||||
func (j *ResizeJob) Run() error {
|
||||
j.mu.RLock()
|
||||
defer j.mu.RUnlock()
|
||||
|
||||
// Set job state to RUNNING.
|
||||
j.setState(ResizeJobStateRunning)
|
||||
j.SetState(ResizeJobStateRunning)
|
||||
|
||||
// Job can be considered done in the case where it doesn't require any action.
|
||||
if !j.urisArePending() {
|
||||
|
|
@ -978,6 +954,10 @@ func (j *ResizeJob) distributeResizeInstructions() error {
|
|||
|
||||
type NodeSet []URI
|
||||
|
||||
func (n NodeSet) Len() int { return len(n) }
|
||||
func (n NodeSet) Swap(i, j int) { n[i], n[j] = n[j], n[i] }
|
||||
func (n NodeSet) Less(i, j int) bool { return n[i].String() < n[j].String() }
|
||||
|
||||
func (u NodeSet) ToHostPortStrings() []string {
|
||||
other := make([]string, 0, len(u))
|
||||
for _, uri := range u {
|
||||
|
|
@ -1012,7 +992,7 @@ func (t *Topology) containsURI(uri URI) bool {
|
|||
return false
|
||||
}
|
||||
|
||||
// AddNode adds the uri to the topology and returns true if added.
|
||||
// AddURI adds the uri to the topology and returns true if added.
|
||||
func (t *Topology) AddURI(uri URI) bool {
|
||||
t.mu.Lock()
|
||||
defer t.mu.Unlock()
|
||||
|
|
@ -1023,6 +1003,11 @@ func (t *Topology) AddURI(uri URI) bool {
|
|||
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.
|
||||
func (c *Cluster) loadTopology() error {
|
||||
buf, err := ioutil.ReadFile(filepath.Join(c.Path, ".topology"))
|
||||
|
|
@ -1117,8 +1102,7 @@ func (c *Cluster) ReceiveEvent(e *NodeEvent) error {
|
|||
return fmt.Errorf("host is not in topology: %v", e.URI)
|
||||
}
|
||||
|
||||
uri := e.URI
|
||||
if err := c.AddNode(uri); err != nil {
|
||||
if err := c.AddNode(e.URI); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
|
|
@ -1138,8 +1122,7 @@ func (c *Cluster) ReceiveEvent(e *NodeEvent) error {
|
|||
|
||||
// If the index does not yet have data, go ahead and add the node.
|
||||
if !c.Holder.HasData() {
|
||||
uri := e.URI
|
||||
if err := c.AddNode(uri); err != nil {
|
||||
if err := c.AddNode(e.URI); err != nil {
|
||||
return err
|
||||
}
|
||||
return c.setStateAndBroadcast(ClusterStateNormal)
|
||||
|
|
@ -1161,7 +1144,7 @@ func (c *Cluster) ReceiveEvent(e *NodeEvent) error {
|
|||
return nil
|
||||
}
|
||||
|
||||
func (c *Cluster) mergeClusterStatus(cs *internal.ClusterStatus) error {
|
||||
func (c *Cluster) MergeClusterStatus(cs *internal.ClusterStatus) error {
|
||||
// Ignore status updates from self (coordinator).
|
||||
if c.IsCoordinator() {
|
||||
return nil
|
||||
|
|
|
|||
248
cluster_test.go
248
cluster_test.go
|
|
@ -15,6 +15,7 @@
|
|||
package pilosa_test
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"math/rand"
|
||||
"reflect"
|
||||
"testing"
|
||||
|
|
@ -213,7 +214,7 @@ func TestCluster_Topology(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
}
|
||||
|
||||
actual := pilosa.Nodes(c1.Nodes).URIs()
|
||||
actual := c1.NodeSet()
|
||||
expected := []pilosa.URI{base, uri1, uri2}
|
||||
|
||||
if !reflect.DeepEqual(actual, expected) {
|
||||
|
|
@ -303,3 +304,248 @@ func TestCluster_Resize(t *testing.T) {
|
|||
}
|
||||
})
|
||||
}
|
||||
|
||||
// TestTestCluster ensures that general cluster functionality works as expected.
|
||||
func TestCluster_ResizeStates(t *testing.T) {
|
||||
|
||||
/* test conditions:
|
||||
x- single node, no data, comes up in NORMAL with topology
|
||||
x- single node, in topology, comes up NORMAL
|
||||
x- single node, not in topology, raises error
|
||||
x- two node, no data, comes up in NORMAL, with topology
|
||||
x- two node, in topology, comes up NORMAL
|
||||
x- two node, STARTING, not in topology, raises error
|
||||
x- two node, NORMAL, not in topology, triggers resize
|
||||
x- resize of nodes with data moves data appropriately
|
||||
*/
|
||||
|
||||
t.Run("Single node, no data", func(t *testing.T) {
|
||||
tc := test.NewTestCluster(1)
|
||||
|
||||
// Open TestCluster.
|
||||
if err := tc.Open(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
node := tc.Clusters[0]
|
||||
|
||||
// Ensure that node comes up in state NORMAL.
|
||||
if node.State != pilosa.ClusterStateNormal {
|
||||
t.Errorf("expected state: %v, but got: %v", pilosa.ClusterStateNormal, node.State)
|
||||
}
|
||||
|
||||
expectedTop := &pilosa.Topology{
|
||||
NodeSet: []pilosa.URI{node.URI},
|
||||
}
|
||||
|
||||
// Verify topology file.
|
||||
if !reflect.DeepEqual(node.Topology, expectedTop) {
|
||||
t.Errorf("expected topology: %v, but got: %v", expectedTop, node.Topology)
|
||||
}
|
||||
|
||||
// Close TestCluster.
|
||||
if err := tc.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Single node, in topology", func(t *testing.T) {
|
||||
tc := test.NewTestCluster(0)
|
||||
tc.AddNode(false)
|
||||
|
||||
node := tc.Clusters[0]
|
||||
|
||||
// write topology to data file
|
||||
top := &pilosa.Topology{
|
||||
NodeSet: []pilosa.URI{node.URI},
|
||||
}
|
||||
tc.WriteTopology(node.Path, top)
|
||||
|
||||
// Open TestCluster.
|
||||
if err := tc.Open(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Ensure that node comes up in state NORMAL.
|
||||
if node.State != pilosa.ClusterStateNormal {
|
||||
t.Errorf("expected state: %v, but got: %v", pilosa.ClusterStateNormal, node.State)
|
||||
}
|
||||
|
||||
// Close TestCluster.
|
||||
if err := tc.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Single node, not in topology", func(t *testing.T) {
|
||||
tc := test.NewTestCluster(0)
|
||||
tc.AddNode(false)
|
||||
|
||||
node := tc.Clusters[0]
|
||||
|
||||
// write topology to data file
|
||||
top := &pilosa.Topology{
|
||||
NodeSet: []pilosa.URI{
|
||||
test.NewURIFromHostPort("some-other-host", 0),
|
||||
},
|
||||
}
|
||||
tc.WriteTopology(node.Path, top)
|
||||
|
||||
// Open TestCluster.
|
||||
expected := "considerTopology: coordinator http://host0:0 is not in topology: [http://some-other-host:0]"
|
||||
err := tc.Open()
|
||||
if err == nil || err.Error() != expected {
|
||||
t.Errorf("did not receive expected error: %s", expected)
|
||||
}
|
||||
|
||||
// Close TestCluster.
|
||||
if err := tc.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Multiple nodes, no data", func(t *testing.T) {
|
||||
tc := test.NewTestCluster(0)
|
||||
tc.AddNode(false)
|
||||
|
||||
// Open TestCluster.
|
||||
if err := tc.Open(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
tc.AddNode(false)
|
||||
|
||||
node0 := tc.Clusters[0]
|
||||
node1 := tc.Clusters[1]
|
||||
|
||||
// Ensure that nodes comes up in state NORMAL.
|
||||
if node0.State != pilosa.ClusterStateNormal {
|
||||
t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateNormal, node0.State)
|
||||
} else if node1.State != pilosa.ClusterStateNormal {
|
||||
t.Errorf("expected node1 state: %v, but got: %v", pilosa.ClusterStateNormal, node1.State)
|
||||
}
|
||||
|
||||
expectedTop := &pilosa.Topology{
|
||||
NodeSet: []pilosa.URI{node0.URI, node1.URI},
|
||||
}
|
||||
|
||||
// Verify topology file.
|
||||
if !reflect.DeepEqual(node0.Topology, expectedTop) {
|
||||
t.Errorf("expected node0 topology: %v, but got: %v", expectedTop, node0.Topology)
|
||||
} else if !reflect.DeepEqual(node1.Topology, expectedTop) {
|
||||
t.Errorf("expected node1 topology: %v, but got: %v", expectedTop, node1.Topology)
|
||||
}
|
||||
|
||||
// Close TestCluster.
|
||||
if err := tc.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Multiple nodes, in/not in topology", func(t *testing.T) {
|
||||
tc := test.NewTestCluster(0)
|
||||
tc.AddNode(false)
|
||||
node0 := tc.Clusters[0]
|
||||
|
||||
u0 := test.NewURIFromHostPort("host0", 0)
|
||||
//u1 := test.NewURIFromHostPort("host1", 0)
|
||||
u2 := test.NewURIFromHostPort("host2", 0)
|
||||
|
||||
// write topology to data file
|
||||
top := &pilosa.Topology{
|
||||
NodeSet: []pilosa.URI{u0, u2},
|
||||
}
|
||||
tc.WriteTopology(node0.Path, top)
|
||||
|
||||
// Open TestCluster.
|
||||
if err := tc.Open(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Ensure that node is in state STARTING before the other node joins.
|
||||
if node0.State != pilosa.ClusterStateStarting {
|
||||
t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateStarting, node0.State)
|
||||
}
|
||||
|
||||
// Expect an error by adding a node not in the topology.
|
||||
expectedError := "host is not in topology: http://host1:0"
|
||||
err := tc.AddNode(false)
|
||||
if err == nil || err.Error() != expectedError {
|
||||
t.Errorf("did not receive expected error: %s", expectedError)
|
||||
}
|
||||
|
||||
tc.AddNode(false)
|
||||
node2 := tc.Clusters[2]
|
||||
|
||||
// Ensure that node comes up in state NORMAL.
|
||||
if node0.State != pilosa.ClusterStateNormal {
|
||||
t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateNormal, node0.State)
|
||||
} else if node2.State != pilosa.ClusterStateNormal {
|
||||
t.Errorf("expected node1 state: %v, but got: %v", pilosa.ClusterStateNormal, node2.State)
|
||||
}
|
||||
|
||||
// Close TestCluster.
|
||||
if err := tc.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("Multiple nodes, with data", func(t *testing.T) {
|
||||
tc := test.NewTestCluster(0)
|
||||
tc.AddNode(false)
|
||||
|
||||
// Open TestCluster.
|
||||
if err := tc.Open(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
|
||||
// Add Data to node0.
|
||||
tc.CreateFrame("i", "f", pilosa.FrameOptions{})
|
||||
tc.SetBit("i", "f", "standard", 1, 101, nil)
|
||||
tc.SetBit("i", "f", "standard", 1, 1300000, nil)
|
||||
|
||||
// AddNode needs to block until the resize process has completed.
|
||||
tc.AddNode(false)
|
||||
|
||||
node0 := tc.Clusters[0]
|
||||
node1 := tc.Clusters[1]
|
||||
|
||||
// Ensure that nodes comes up in state NORMAL.
|
||||
if node0.State != pilosa.ClusterStateNormal {
|
||||
t.Errorf("expected node0 state: %v, but got: %v", pilosa.ClusterStateNormal, node0.State)
|
||||
} else if node1.State != pilosa.ClusterStateNormal {
|
||||
t.Errorf("expected node1 state: %v, but got: %v", pilosa.ClusterStateNormal, node1.State)
|
||||
}
|
||||
|
||||
expectedTop := &pilosa.Topology{
|
||||
NodeSet: []pilosa.URI{node0.URI, node1.URI},
|
||||
}
|
||||
|
||||
// Verify topology file.
|
||||
if !reflect.DeepEqual(node0.Topology, expectedTop) {
|
||||
t.Errorf("expected node0 topology: %v, but got: %v", expectedTop, node0.Topology)
|
||||
} else if !reflect.DeepEqual(node1.Topology, expectedTop) {
|
||||
t.Errorf("expected node1 topology: %v, but got: %v", expectedTop, node1.Topology)
|
||||
}
|
||||
|
||||
// Verify that node-1 contains the fragment (i/f/standard/1) transferred from node-0.
|
||||
node0Frame := node0.Holder.Frame("i", "f")
|
||||
node0View := node0Frame.View("standard")
|
||||
node0Fragment := node0View.Fragment(1)
|
||||
|
||||
node1Frame := node1.Holder.Frame("i", "f")
|
||||
node1View := node1Frame.View("standard")
|
||||
node1Fragment := node1View.Fragment(1)
|
||||
|
||||
// Ensure checksums are the same.
|
||||
orig := node0Fragment.Checksum()
|
||||
if chksum := node1Fragment.Checksum(); !bytes.Equal(chksum, orig) {
|
||||
t.Fatalf("expected checksum to match: %x - %x", chksum, orig)
|
||||
}
|
||||
|
||||
// Close TestCluster.
|
||||
if err := tc.Close(); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
|
|
|||
|
|
@ -781,7 +781,7 @@ func TestExecutor_Execute_FieldRange(t *testing.T) {
|
|||
t.Fatal(err)
|
||||
} else if !reflect.DeepEqual([]uint64{0}, result[0].(*pilosa.Bitmap).Bits()) {
|
||||
//t.Fatalf("unexpected result: %s", spew.Sdump(result))
|
||||
t.Fatalf("unexpected result: %s", result[0].(*pilosa.Bitmap).Bits())
|
||||
t.Fatalf("unexpected result: %v", result[0].(*pilosa.Bitmap).Bits())
|
||||
}
|
||||
})
|
||||
|
||||
|
|
|
|||
17
frame.go
17
frame.go
|
|
@ -550,6 +550,18 @@ func (f *Frame) Views() []*View {
|
|||
return other
|
||||
}
|
||||
|
||||
// viewNames returns a list of all views (as a string) in the frame.
|
||||
func (f *Frame) viewNames() []string {
|
||||
f.mu.Lock()
|
||||
defer f.mu.Unlock()
|
||||
|
||||
other := make([]string, 0, len(f.views))
|
||||
for viewName, _ := range f.views {
|
||||
other = append(other, viewName)
|
||||
}
|
||||
return other
|
||||
}
|
||||
|
||||
// RecalculateCaches recalculates caches on every view in the frame.
|
||||
func (f *Frame) RecalculateCaches() {
|
||||
for _, view := range f.Views() {
|
||||
|
|
@ -958,8 +970,9 @@ func encodeFrames(a []*Frame) []*internal.Frame {
|
|||
func encodeFrame(f *Frame) *internal.Frame {
|
||||
fo := f.options()
|
||||
return &internal.Frame{
|
||||
Name: f.name,
|
||||
Meta: fo.Encode(),
|
||||
Name: f.name,
|
||||
Meta: fo.Encode(),
|
||||
Views: f.viewNames(),
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
29
holder.go
29
holder.go
|
|
@ -187,6 +187,35 @@ func (h *Holder) Schema() []*IndexInfo {
|
|||
return a
|
||||
}
|
||||
|
||||
// ApplySchema applies an internal Schema to Holder.
|
||||
func (h *Holder) ApplySchema(schema *internal.Schema) error {
|
||||
// Create indexes that don't exist.
|
||||
for _, index := range schema.Indexes {
|
||||
opt := IndexOptions{}
|
||||
idx, err := h.CreateIndexIfNotExists(index.Name, opt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// Create frames that don't exist.
|
||||
for _, f := range index.Frames {
|
||||
opt := decodeFrameOptions(f.Meta)
|
||||
frame, err := idx.CreateFrameIfNotExists(f.Name, *opt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// Create views that don't exist.
|
||||
for _, v := range f.Views {
|
||||
_, err := frame.CreateViewIfNotExists(v)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
}
|
||||
// TODO: Create inputDefinitions that don't exist.
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// EncodeMaxSlices creates and internal representation of max slices.
|
||||
func (h *Holder) EncodeMaxSlices() *internal.MaxSlices {
|
||||
return &internal.MaxSlices{
|
||||
|
|
|
|||
25
server.go
25
server.go
|
|
@ -332,12 +332,15 @@ func (s *Server) ReceiveMessage(pb proto.Message) error {
|
|||
return err
|
||||
}
|
||||
case *internal.ClusterStatus:
|
||||
err := s.Cluster.mergeClusterStatus(obj)
|
||||
err := s.Cluster.MergeClusterStatus(obj)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.ResizeInstruction:
|
||||
s.Cluster.followResizeInstruction(obj)
|
||||
err := s.Cluster.FollowResizeInstruction(obj)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
case *internal.ResizeInstructionComplete:
|
||||
err := s.Cluster.MarkResizeInstructionComplete(obj)
|
||||
if err != nil {
|
||||
|
|
@ -397,22 +400,8 @@ func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error {
|
|||
}
|
||||
|
||||
// Sync schema.
|
||||
// Create indexes that don't exist.
|
||||
for _, index := range ns.Schema.Indexes {
|
||||
opt := IndexOptions{}
|
||||
idx, err := s.Holder.CreateIndexIfNotExists(index.Name, opt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// Create frames that don't exist.
|
||||
for _, f := range index.Frames {
|
||||
opt := decodeFrameOptions(f.Meta)
|
||||
_, err := idx.CreateFrameIfNotExists(f.Name, *opt)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
// TODO: Create inputDefinitions that don't exist.
|
||||
if err := s.Holder.ApplySchema(ns.Schema); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Sync maxSlices (standard).
|
||||
|
|
|
|||
307
test/cluster.go
307
test/cluster.go
|
|
@ -1,10 +1,17 @@
|
|||
package test
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"bytes"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"path/filepath"
|
||||
"sort"
|
||||
"time"
|
||||
|
||||
"github.com/gogo/protobuf/proto"
|
||||
"github.com/pilosa/pilosa"
|
||||
"github.com/pilosa/pilosa/internal"
|
||||
)
|
||||
|
||||
// NewCluster returns a cluster with n nodes and uses a mod-based hasher.
|
||||
|
|
@ -62,3 +69,303 @@ func NewURIFromHostPort(host string, port uint16) pilosa.URI {
|
|||
uri.SetPort(port)
|
||||
return *uri
|
||||
}
|
||||
|
||||
// TestCluster represents a cluster of test nodes, each of which
|
||||
// has a pilosa.Cluster.
|
||||
type TestCluster struct {
|
||||
Clusters []*pilosa.Cluster
|
||||
|
||||
common *commonClusterSettings
|
||||
|
||||
resizeDone chan struct{}
|
||||
}
|
||||
|
||||
type commonClusterSettings struct {
|
||||
NodeSet pilosa.NodeSet
|
||||
}
|
||||
|
||||
func (t *TestCluster) CreateIndex(name string) error {
|
||||
for _, c := range t.Clusters {
|
||||
if _, err := c.Holder.CreateIndexIfNotExists(name, pilosa.IndexOptions{}); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (t *TestCluster) CreateFrame(index, frame string, opt pilosa.FrameOptions) error {
|
||||
for _, c := range t.Clusters {
|
||||
idx, err := c.Holder.CreateIndexIfNotExists(index, pilosa.IndexOptions{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := idx.CreateFrame(frame, opt); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
func (t *TestCluster) SetBit(index, frame, view 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 slice location.
|
||||
slice := colID / pilosa.SliceWidth
|
||||
nodes := c0.FragmentNodes(index, slice)
|
||||
|
||||
for _, node := range nodes {
|
||||
c := t.clusterByURI(node.URI)
|
||||
if c == nil {
|
||||
continue
|
||||
}
|
||||
f := c.Holder.Frame(index, frame)
|
||||
if f == nil {
|
||||
return fmt.Errorf("index/frame does not exist: %s/%s", index, frame)
|
||||
}
|
||||
_, err := f.SetBit(view, rowID, colID, x)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (t *TestCluster) clusterByURI(uri pilosa.URI) *pilosa.Cluster {
|
||||
for _, c := range t.Clusters {
|
||||
if c.URI == uri {
|
||||
return c
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// AddNode adds a node to the cluster and (potentially) starts a resize job.
|
||||
func (t *TestCluster) AddNode(saveTopology bool) error {
|
||||
id := len(t.Clusters)
|
||||
|
||||
c, err := t.addCluster(id, saveTopology)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Send NodeJoin event to coordinator.
|
||||
if id > 0 {
|
||||
coord := t.Clusters[0]
|
||||
ev := &pilosa.NodeEvent{
|
||||
Event: pilosa.NodeJoin,
|
||||
URI: c.URI,
|
||||
}
|
||||
|
||||
//go coord.ReceiveEvent(ev)
|
||||
if err := coord.ReceiveEvent(ev); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Wait for the AddNode job to finish.
|
||||
if c.State != pilosa.ClusterStateNormal {
|
||||
t.resizeDone = make(chan struct{})
|
||||
<-t.resizeDone
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// WriteTopology writes the given topology to disk.
|
||||
func (t *TestCluster) WriteTopology(path string, top *pilosa.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 *TestCluster) addCluster(i int, saveTopology bool) (*pilosa.Cluster, error) {
|
||||
|
||||
uri := NewURI("http", fmt.Sprintf("host%d", i), uint16(0))
|
||||
|
||||
// add URI to common
|
||||
t.common.NodeSet = append(t.common.NodeSet, uri)
|
||||
sort.Sort(t.common.NodeSet)
|
||||
|
||||
// create node-specific temp directory
|
||||
path, err := ioutil.TempDir("", fmt.Sprintf("pilosa-cluster-node-%d-", i))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// holder
|
||||
h := pilosa.NewHolder()
|
||||
h.Path = path
|
||||
|
||||
// cluster
|
||||
c := pilosa.NewCluster()
|
||||
c.ReplicaN = 1
|
||||
c.Hasher = NewModHasher()
|
||||
c.Path = path
|
||||
c.Topology = pilosa.NewTopology()
|
||||
c.Holder = h
|
||||
c.MemberSet = pilosa.NewStaticMemberSet()
|
||||
c.URI = uri
|
||||
c.Coordinator = t.common.NodeSet[0] // the first node is the coordinator
|
||||
c.Broadcaster = t
|
||||
|
||||
// add nodes
|
||||
if saveTopology {
|
||||
for _, u := range t.common.NodeSet {
|
||||
c.AddNode(u)
|
||||
}
|
||||
}
|
||||
|
||||
// Add this node to the TestCluster.
|
||||
t.Clusters = append(t.Clusters, c)
|
||||
|
||||
return c, nil
|
||||
}
|
||||
|
||||
// NewTestCluster returns a new instance of test.Cluster.
|
||||
func NewTestCluster(n int) *TestCluster {
|
||||
|
||||
tc := &TestCluster{
|
||||
common: &commonClusterSettings{},
|
||||
}
|
||||
|
||||
// add clusters
|
||||
for i := 0; i < n; i++ {
|
||||
_, err := tc.addCluster(i, true)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
}
|
||||
return tc
|
||||
}
|
||||
|
||||
// SetState sets the state of the cluster on each node.
|
||||
func (t *TestCluster) SetState(state string) {
|
||||
for _, c := range t.Clusters {
|
||||
c.State = state
|
||||
}
|
||||
}
|
||||
|
||||
// Open opens all clusters in the test cluster.
|
||||
func (t *TestCluster) Open() error {
|
||||
for _, c := range t.Clusters {
|
||||
err := c.Open()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// Close closes all clusters in the test cluster.
|
||||
func (t *TestCluster) Close() error {
|
||||
for _, c := range t.Clusters {
|
||||
err := c.Close()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// TestCluster implements Broadcaster interface.
|
||||
|
||||
// SendSync is a test implemenetation of Broadcaster SendSync method.
|
||||
func (t *TestCluster) SendSync(pb proto.Message) error {
|
||||
switch obj := pb.(type) {
|
||||
case *internal.ClusterStatus:
|
||||
// Apply the send message to all nodes (except the coordinator).
|
||||
for _, c := range t.Clusters {
|
||||
c.MergeClusterStatus(obj)
|
||||
}
|
||||
if obj.State == pilosa.ClusterStateNormal && t.resizeDone != nil {
|
||||
close(t.resizeDone)
|
||||
}
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// SendAsync is a test implemenetation of Broadcaster SendAsync method.
|
||||
func (t *TestCluster) SendAsync(pb proto.Message) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// SendTo is a test implemenetation of Broadcaster SendTo method.
|
||||
func (t *TestCluster) SendTo(to *pilosa.Node, pb proto.Message) error {
|
||||
switch obj := pb.(type) {
|
||||
case *internal.ResizeInstruction:
|
||||
t.FollowResizeInstruction(obj)
|
||||
case *internal.ResizeInstructionComplete:
|
||||
coord := t.clusterByURI(to.URI)
|
||||
go coord.MarkResizeInstructionComplete(obj)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// FollowResizeInstruction is a version of cluster.FollowResizeInstruction used for testing.
|
||||
func (t *TestCluster) FollowResizeInstruction(instr *internal.ResizeInstruction) error {
|
||||
|
||||
// Prepare the return message.
|
||||
complete := &internal.ResizeInstructionComplete{
|
||||
JobID: instr.JobID,
|
||||
URI: instr.URI,
|
||||
Error: "",
|
||||
}
|
||||
|
||||
// figure out which node it was meant for, then call the operation on that cluster
|
||||
// basically need to mimic this: client.RetrieveSliceFromURI(context.Background(), src.Index, src.Frame, src.View, src.Slice, srcURI)
|
||||
instrURI := pilosa.DecodeURI(instr.URI)
|
||||
destCluster := t.clusterByURI(instrURI)
|
||||
|
||||
// Sync the schema received in the resize instruction.
|
||||
if err := destCluster.Holder.ApplySchema(instr.Schema); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
for _, src := range instr.Sources {
|
||||
srcURI := pilosa.DecodeURI(src.URI)
|
||||
srcCluster := t.clusterByURI(srcURI)
|
||||
|
||||
srcFragment := srcCluster.Holder.Fragment(src.Index, src.Frame, src.View, src.Slice)
|
||||
destFragment := destCluster.Holder.Fragment(src.Index, src.Frame, src.View, src.Slice)
|
||||
if destFragment == nil {
|
||||
// Create fragment on destination if it doesn't exist.
|
||||
f := destCluster.Holder.Frame(src.Index, src.Frame)
|
||||
v := f.View(src.View)
|
||||
var err error
|
||||
destFragment, err = v.CreateFragmentIfNotExists(src.Slice)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
buf := bytes.NewBuffer(nil)
|
||||
|
||||
bw := bufio.NewWriter(buf)
|
||||
br := bufio.NewReader(buf)
|
||||
|
||||
// Get the fragment from source.
|
||||
if _, err := srcFragment.WriteTo(bw); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Flush the bufio.buf to the io.Writer (buf).
|
||||
bw.Flush()
|
||||
|
||||
// Write data to destination.
|
||||
if _, err := destFragment.ReadFrom(br); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
node := &pilosa.Node{
|
||||
URI: pilosa.DecodeURI(instr.Coordinator),
|
||||
}
|
||||
if err := t.SendTo(node, complete); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
|
|
|||
4
uri.go
4
uri.go
|
|
@ -204,6 +204,10 @@ func encodeURI(u URI) *internal.URI {
|
|||
}
|
||||
}
|
||||
|
||||
func DecodeURI(i *internal.URI) URI {
|
||||
return decodeURI(i)
|
||||
}
|
||||
|
||||
func decodeURI(i *internal.URI) URI {
|
||||
if i == nil {
|
||||
return URI{}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue