diff --git a/client_test.go b/client_test.go index e5dffc9a5..15ebb8996 100644 --- a/client_test.go +++ b/client_test.go @@ -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) } } diff --git a/cluster.go b/cluster.go index 2841bd5bc..323283a8a 100644 --- a/cluster.go +++ b/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 diff --git a/cluster_test.go b/cluster_test.go index 6fa3ed6e5..a51823993 100644 --- a/cluster_test.go +++ b/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) + } + }) +} diff --git a/executor_test.go b/executor_test.go index 34fe21e9a..040248de9 100644 --- a/executor_test.go +++ b/executor_test.go @@ -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()) } }) diff --git a/frame.go b/frame.go index 06d081423..c57edf38e 100644 --- a/frame.go +++ b/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(), } } diff --git a/holder.go b/holder.go index 8da92f498..ce6e052e2 100644 --- a/holder.go +++ b/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{ diff --git a/server.go b/server.go index 9327e8ca8..e460e1432 100644 --- a/server.go +++ b/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). diff --git a/test/cluster.go b/test/cluster.go index 604aff14c..e73289827 100644 --- a/test/cluster.go +++ b/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 +} diff --git a/uri.go b/uri.go index 23aba4dfc..79c7cd4ac 100644 --- a/uri.go +++ b/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{}