diff --git a/cluster.go b/cluster.go index fa6651eab..befc66c1a 100644 --- a/cluster.go +++ b/cluster.go @@ -51,6 +51,9 @@ const ( // Final states. ResizeJobStateDone = "DONE" ResizeJobStateAborted = "ABORTED" + + ResizeJobActionAdd = "ADD" + ResizeJobActionRemove = "REMOVE" ) // Node represents a node in the cluster. @@ -127,6 +130,12 @@ func (h ByHost) Len() int { return len(h) } func (h ByHost) Swap(i, j int) { h[i], h[j] = h[j], h[i] } func (h ByHost) Less(i, j int) bool { return h[i].URI.String() < h[j].URI.String() } +// nodeAction represents a node that is joining or leaving the cluster. +type nodeAction struct { + uri URI + action string +} + // Cluster represents a collection of nodes. type Cluster struct { URI URI @@ -158,7 +167,7 @@ type Cluster struct { Holder *Holder Broadcaster Broadcaster - joiningURIs chan URI + joiningLeavingNodes chan nodeAction mu sync.RWMutex jobs map[int64]*ResizeJob @@ -181,9 +190,9 @@ func NewCluster() *Cluster { ReplicaN: DefaultReplicaN, EventReceiver: NopEventReceiver, - joiningURIs: make(chan URI, 10), // buffered channel - jobs: make(map[int64]*ResizeJob), - closing: make(chan struct{}), + joiningLeavingNodes: make(chan nodeAction, 10), // buffered channel + jobs: make(map[int64]*ResizeJob), + closing: make(chan struct{}), LogOutput: os.Stderr, prefect: &NopSecurityManager{}, @@ -215,7 +224,7 @@ func (c *Cluster) SetCoordinator(oldURI, newURI URI) bool { func (c *Cluster) AddNode(uri URI) error { // add to cluster - _, added := c.AddNodeBasicSorted(uri) + _, added := c.addNodeBasicSorted(uri) if !added { return nil } @@ -232,6 +241,27 @@ func (c *Cluster) AddNode(uri URI) error { return c.saveTopology() } +// RemoveNode removes a node from the Cluster and updates and saves the +// new topology. +func (c *Cluster) RemoveNode(uri URI) error { + // remove from cluster + removed := c.removeNodeBasicSorted(uri) + if !removed { + return nil + } + + // remove from topology + if c.Topology == nil { + return fmt.Errorf("Cluster.Topology is nil") + } + if !c.Topology.RemoveURI(uri) { + return nil + } + + // save topology + return c.saveTopology() +} + // NodeSet returns the list of uris in the cluster. func (c *Cluster) NodeSet() []URI { return Nodes(c.Nodes).URIs() @@ -278,9 +308,19 @@ func (c *Cluster) NodeByURI(uri URI) *Node { return nil } -// AddNodeBasicSorted adds a node to the cluster, sorted by uri. +// nodePositionByURI returns the position of the node in slice c.Nodes. +func (c *Cluster) nodePositionByURI(uri URI) int { + for i, n := range c.Nodes { + if n.URI == uri { + return i + } + } + return -1 +} + +// addNodeBasicSorted adds a node to the cluster, sorted by uri. // Returns a pointer to the node and true if the node was added. -func (c *Cluster) AddNodeBasicSorted(uri URI) (*Node, bool) { +func (c *Cluster) addNodeBasicSorted(uri URI) (*Node, bool) { n := c.NodeByURI(uri) if n != nil { return n, false @@ -295,6 +335,21 @@ func (c *Cluster) AddNodeBasicSorted(uri URI) (*Node, bool) { return n, true } +// removeNodeBasicSorted removes a node from the cluster, maintaining +// the sort order. Returns true if the node was removed. +func (c *Cluster) removeNodeBasicSorted(uri URI) bool { + i := c.nodePositionByURI(uri) + if i < 0 { + return false + } + + copy(c.Nodes[i:], c.Nodes[i+1:]) + c.Nodes[len(c.Nodes)-1] = nil + c.Nodes = c.Nodes[:len(c.Nodes)-1] + + return true +} + // frag is a struct of basic fragment information. type frag struct { frame string @@ -363,8 +418,8 @@ func (c *Cluster) fragsByHost(idx *Index) fragsByHost { func (c *Cluster) fragCombos(idx string, maxSlice uint64, frameViews viewsByFrame) fragsByHost { t := make(fragsByHost) for i := uint64(0); i <= maxSlice; i++ { - f := c.FragmentNodes(idx, i) - for _, n := range f { + nodes := c.FragmentNodes(idx, i) + for _, n := range nodes { // for each frame/view combination: for frame, views := range frameViews { for _, view := range views { @@ -376,21 +431,70 @@ func (c *Cluster) fragCombos(idx string, maxSlice uint64, frameViews viewsByFram return t } -// DataDiff returns a list of ResizeSources - for each host in the `to` cluster - +// diff compares c with another cluster and determines if a node is being +// added or removed. An error is returned for any case other than where +// exactly one node is added or removed. +func (c *Cluster) diff(other *Cluster) (action string, uri URI, err error) { + lenFrom := len(c.Nodes) + lenTo := len(other.Nodes) + // Determine if a node is being added or removed. + if lenFrom == lenTo { + return action, uri, errors.New("clusters are the same size") + } + if lenFrom < lenTo { + // Adding a node. + if lenTo-lenFrom > 1 { + return action, uri, errors.New("adding more than one node at a time is not supported") + } + action = ResizeJobActionAdd + // Determine the URI that is being added. + for _, n := range other.Nodes { + if c.NodeByURI(n.URI) == nil { + uri = n.URI + break + } + } + } else if len(c.Nodes) > len(other.Nodes) { + // Removing a node. + if lenFrom-lenTo > 1 { + return action, uri, errors.New("removing more than one node at a time is not supported") + } + action = ResizeJobActionRemove + // Determine the URI that is being removed. + for _, n := range c.Nodes { + if other.NodeByURI(n.URI) == nil { + uri = n.URI + break + } + } + } + return action, uri, nil +} + +// fragSources returns a list of ResizeSources - for each node in the `to` cluster - // required to move from cluster `c` to cluster `to`. -func (c *Cluster) DataDiff(to *Cluster, idx *Index) map[URI][]*internal.ResizeSource { +func (c *Cluster) fragSources(to *Cluster, idx *Index) (map[URI][]*internal.ResizeSource, error) { m := make(map[URI][]*internal.ResizeSource) + // Determine if a node is being added or removed. + action, diffURI, err := c.diff(to) + if err != nil { + return nil, err + } + // Initialize the map with all the nodes in `to`. for _, n := range to.Nodes { m[n.URI] = nil } - // For now, we want our source to be confined to the primary fragment - // (i.e. don't use replicas as source data). So if it's not already, - // base our source fragments on a cluster with replica = 1. + // If a node is being added, the source can be confined to the + // primary fragments (i.e. no need to use replicas as source data). + // In this case, source fragments can be based on a cluster with + // replica = 1. + // If a node is being removed, however, then it will most likely + // require that a replica fragment be the source data. srcCluster := c - if c.ReplicaN > 1 { + if action == ResizeJobActionAdd && c.ReplicaN > 1 { srcCluster = NewCluster() srcCluster.Nodes = Nodes(c.Nodes).Clone() srcCluster.Hasher = c.Hasher @@ -408,6 +512,10 @@ func (c *Cluster) DataDiff(to *Cluster, idx *Index) map[URI][]*internal.ResizeSo // srcHostsByFrag is the inverse representation of srcFrags. srcHostsByFrag := make(map[frag]URI) for uri, frags := range srcFrags { + // If a node is being removed, don't consider it as a source. + if action == ResizeJobActionRemove && uri == diffURI { + continue + } for _, frag := range frags { srcHostsByFrag[frag] = uri } @@ -427,18 +535,28 @@ func (c *Cluster) DataDiff(to *Cluster, idx *Index) map[URI][]*internal.ResizeSo for host, diff := range diffs { m[host] = []*internal.ResizeSource{} for _, frag := range diff { + // If there is no valid source URI for a fragment, + // it likely means that the replica factor was not + // high enough for the remaining nodes to contain + // the fragment. + srcHost, ok := srcHostsByFrag[frag] + if !ok { + return nil, errors.New("not enough data to perform resize") + } + src := &internal.ResizeSource{ - URI: (srcHostsByFrag[frag]).Encode(), + URI: (srcHost).Encode(), Index: idx.Name(), Frame: frag.frame, View: frag.view, Slice: frag.slice, } + m[host] = append(m[host], src) } } - return m + return m, nil } // Partition returns the partition that a slice belongs to. @@ -575,8 +693,8 @@ func (c *Cluster) haveTopologyAgreement() bool { return URISlicesAreEqual(c.Topology.NodeSet, c.NodeSet()) } -func (c *Cluster) handleJoiningHost(uri URI) error { - j, err := c.GenerateResizeJob(uri) +func (c *Cluster) handleNodeAction(nodeAction nodeAction) error { + j, err := c.generateResizeJob(nodeAction) if err != nil { return err } @@ -594,8 +712,12 @@ func (c *Cluster) handleJoiningHost(uri URI) error { if err := c.CompleteCurrentJob(ResizeJobStateDone); err != nil { return err } - // Add uri to the cluster. - return c.AddNode(uri) + // Add/remove uri to/from the cluster. + if j.action == ResizeJobActionRemove { + return c.RemoveNode(nodeAction.uri) + } else if j.action == ResizeJobActionAdd { + return c.AddNode(nodeAction.uri) + } case ResizeJobStateAborted: if err := c.CompleteCurrentJob(ResizeJobStateAborted); err != nil { return err @@ -623,10 +745,10 @@ func (c *Cluster) listenForJoins() { // Handle all pending joins before changing state back to NORMAL. select { - case uri := <-c.joiningURIs: - err := c.handleJoiningHost(uri) + case nodeAction := <-c.joiningLeavingNodes: + err := c.handleNodeAction(nodeAction) if err != nil { - c.logger().Printf("handleJoiningHost error: err=%s", err) + c.logger().Printf("handleNodeAction error: err=%s", err) continue } uriJoined = true @@ -646,10 +768,10 @@ func (c *Cluster) listenForJoins() { select { case <-c.closing: return - case uri := <-c.joiningURIs: - err := c.handleJoiningHost(uri) + case nodeAction := <-c.joiningLeavingNodes: + err := c.handleNodeAction(nodeAction) if err != nil { - c.logger().Printf("handleJoiningHost error: err=%s", err) + c.logger().Printf("handleNodeAction error: err=%s", err) continue } uriJoined = true @@ -658,14 +780,17 @@ func (c *Cluster) listenForJoins() { } } -// GenerateResizeJob creates a new ResizeJob based on the new host being -// added. It also saves a reference to the ResizeJob in the `jobs` map +// generateResizeJob creates a new ResizeJob based on the new node being +// added/removed. It also saves a reference to the ResizeJob in the `jobs` map // for future lookup by JobID. -func (c *Cluster) GenerateResizeJob(addURI URI) (*ResizeJob, error) { +func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*ResizeJob, error) { c.mu.Lock() defer c.mu.Unlock() - j := c.generateResizeJob(addURI) + j, err := c.generateResizeJobByAction(nodeAction) + if err != nil { + return nil, err + } // Save job in jobs map for future reference. c.jobs[j.ID] = j @@ -679,51 +804,57 @@ func (c *Cluster) GenerateResizeJob(addURI URI) (*ResizeJob, error) { return j, nil } -// generateResizeJob returns a ResizeJob with instructions based on -// the difference between Cluster and a new Cluster containing addHost. +// generateResizeJobByAction returns a ResizeJob with instructions based on +// the difference between Cluster and a new Cluster with/without uri. // Broadcaster is associated to the ResizeJob here for use in broadcasting // the resize instructions to other nodes in the cluster. -func (c *Cluster) generateResizeJob(addURI URI) *ResizeJob { +func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*ResizeJob, error) { - j := NewResizeJob(addURI, Nodes(c.Nodes).URIs()) + j := NewResizeJob(Nodes(c.Nodes).URIs(), nodeAction.uri, nodeAction.action) j.Broadcaster = c.Broadcaster - // toCluster is a clone of Cluster with the new node added for comparison. + // toCluster is a clone of Cluster with the new node added/removed for comparison. toCluster := NewCluster() toCluster.Nodes = Nodes(c.Nodes).Clone() toCluster.Hasher = c.Hasher toCluster.PartitionN = c.PartitionN toCluster.ReplicaN = c.ReplicaN - toCluster.AddNodeBasicSorted(addURI) + if nodeAction.action == ResizeJobActionRemove { + toCluster.removeNodeBasicSorted(nodeAction.uri) + } else if nodeAction.action == ResizeJobActionAdd { + toCluster.addNodeBasicSorted(nodeAction.uri) + } pbSchema := c.Holder.EncodeSchema() // Add to the ResizeJob the instructions for each index. for _, idx := range c.Holder.Indexes() { - // dataDiff is map[string][]*internal.ResizeSource, where string is - // a host in toCluster. - dataDiff := c.DataDiff(toCluster, idx) + // fragSources is map[URI][]*internal.ResizeSource. + fragSources, err := c.fragSources(toCluster, idx) + if err != nil { + return nil, err + } - for uri, sources := range dataDiff { + for u, sources := range fragSources { // If a host doesn't need to request data, mark it as complete. if len(sources) == 0 { - j.URIs[uri] = true + j.URIs[u] = true continue } - // TODO: we can probably consilidate the instructions that go to the same + // TODO: we can probably consolidate the instructions that go to the same // node but apply to different indexes. (i.e. don't nest this in the Indexes() loop) instr := &internal.ResizeInstruction{ JobID: j.ID, - URI: uri.Encode(), + URI: u.Encode(), Coordinator: encodeURI(c.Coordinator), Sources: sources, - Schema: pbSchema, + Schema: pbSchema, // Include the schema to ensure it's in sync on the receiving node. } j.Instructions = append(j.Instructions, instr) } } - return j + return j, nil } // CompleteCurrentJob sets the state of the current ResizeJob @@ -862,6 +993,7 @@ type ResizeJob struct { Instructions []*internal.ResizeInstruction Broadcaster Broadcaster + action string result chan string mu sync.RWMutex @@ -869,22 +1001,33 @@ type ResizeJob struct { } // NewResizeJob returns a new instance of ResizeJob. -func NewResizeJob(addURI URI, existingURIs []URI) *ResizeJob { +func NewResizeJob(existingURIs []URI, uri URI, action string) *ResizeJob { // Build a map of uris to track their resize status. - uris := make(map[URI]bool) - // The value for a node will be set to true after that node // has indicated that it has completed all resize instructions. - for _, u := range existingURIs { - uris[u] = false + uris := make(map[URI]bool) + + if action == ResizeJobActionRemove { + for _, u := range existingURIs { + // Exclude the removed node from the map. + if u == uri { + continue + } + uris[u] = false + } + } else if action == ResizeJobActionAdd { + for _, u := range existingURIs { + uris[u] = false + } + // Include the added node in the map for tracking. + uris[uri] = false } - // Include the added node in the map for tracking. - uris[addURI] = false return &ResizeJob{ ID: rand.Int63(), URIs: uris, + action: action, result: make(chan string), } } @@ -1009,6 +1152,15 @@ func (t *Topology) containsURI(uri URI) bool { return false } +func (t *Topology) positionByURI(uri URI) int { + for i, turi := range t.NodeSet { + if turi == uri { + return i + } + } + return -1 +} + // AddURI adds the uri to the topology and returns true if added. func (t *Topology) AddURI(uri URI) bool { t.mu.Lock() @@ -1020,6 +1172,23 @@ func (t *Topology) AddURI(uri URI) bool { return true } +// RemoveURI removes the uri from the topology and returns true if removed. +func (t *Topology) RemoveURI(uri URI) bool { + t.mu.Lock() + defer t.mu.Unlock() + + i := t.positionByURI(uri) + if i < 0 { + return false + } + + copy(t.NodeSet[i:], t.NodeSet[i+1:]) + t.NodeSet[len(t.NodeSet)-1] = URI{} + t.NodeSet = t.NodeSet[:len(t.NodeSet)-1] + + return true +} + // Encode converts t into its internal representation. func (t *Topology) Encode() *internal.Topology { return encodeTopology(t) @@ -1112,52 +1281,94 @@ func (c *Cluster) ReceiveEvent(e *NodeEvent) error { if !c.IsCoordinator() { return nil } + return c.nodeJoin(e.URI) + case NodeLeave: + // Automatic nodeLeave is intentionally not implemented. + case NodeUpdate: + // NodeUpdate is intentionally not implemented. + } - if c.needTopologyAgreement() { - // A host that is not part of the topology can't be added to the STARTING cluster. - if !c.Topology.ContainsURI(e.URI) { - return fmt.Errorf("host is not in topology: %v", e.URI) - } + return nil +} - if err := c.AddNode(e.URI); err != nil { - return err - } - - // If the result of the previous AddNode completed the joining of nodes - // in the topology, then change the state to NORMAL. - if c.haveTopologyAgreement() { - return c.setStateAndBroadcast(ClusterStateNormal) - } - - return nil +func (c *Cluster) nodeJoin(uri URI) error { + if c.needTopologyAgreement() { + // A host that is not part of the topology can't be added to the STARTING cluster. + if !c.Topology.ContainsURI(uri) { + return fmt.Errorf("host is not in topology: %v", uri) } - // Don't do anything else if the cluster already contains the node. - if c.NodeByURI(e.URI) != nil { - return nil + if err := c.AddNode(uri); err != nil { + return err } - // If the index does not yet have data, go ahead and add the node. - if !c.Holder.HasData() { - if err := c.AddNode(e.URI); err != nil { - return err - } + // If the result of the previous AddNode completed the joining of nodes + // in the topology, then change the state to NORMAL. + if c.haveTopologyAgreement() { return c.setStateAndBroadcast(ClusterStateNormal) } - // If the cluster has data, we need to change to RESIZING and - // kick off the resizing process. - if err := c.setStateAndBroadcast(ClusterStateResizing); err != nil { + return nil + } + + // Don't do anything else if the cluster already contains the node. + if c.NodeByURI(uri) != nil { + return nil + } + + // If the holder does not yet contain data, go ahead and add the node. + if !c.Holder.HasData() { + if err := c.AddNode(uri); err != nil { return err } - c.joiningURIs <- e.URI - - case NodeLeave: - // TODO: implement this - case NodeUpdate: - // TODO: implement this + return c.setStateAndBroadcast(ClusterStateNormal) } + // If the cluster has data, we need to change to RESIZING and + // kick off the resizing process. + if err := c.setStateAndBroadcast(ClusterStateResizing); err != nil { + return err + } + c.joiningLeavingNodes <- nodeAction{uri, ResizeJobActionAdd} + + return nil +} + +// NodeLeave initiates the removal of a node from the cluster. +func (c *Cluster) NodeLeave(uri URI) error { + // Refuse the request if this is not the coordinator. + if !c.IsCoordinator() { + return fmt.Errorf("Node removal requests are only valid on the Coordinator node: %s", c.Coordinator) + } + + if c.State != ClusterStateNormal { + return fmt.Errorf("Cluster must be in state %s to remove a node. Current state: %s", ClusterStateNormal, c.State) + } + + return c.nodeLeave(uri) +} + +func (c *Cluster) nodeLeave(uri URI) error { + // Don't do anything else if the cluster doesn't contain the node. + if c.NodeByURI(uri) == nil { + return nil + } + + // If the holder does not yet contain data, go ahead and remove the node. + if !c.Holder.HasData() { + if err := c.RemoveNode(uri); err != nil { + return err + } + return c.setStateAndBroadcast(ClusterStateNormal) + } + + // If the cluster has data then change state to RESIZING and + // kick off the resizing process. + if err := c.setStateAndBroadcast(ClusterStateResizing); err != nil { + return err + } + c.joiningLeavingNodes <- nodeAction{uri, ResizeJobActionRemove} + return nil } diff --git a/cluster_internal_test.go b/cluster_internal_test.go index 4ef65c958..75eb86f85 100644 --- a/cluster_internal_test.go +++ b/cluster_internal_test.go @@ -15,14 +15,16 @@ package pilosa import ( + "io/ioutil" "reflect" "testing" + + "github.com/pilosa/pilosa/internal" ) // Ensure that fragCombos creates the correct fragment mapping. func TestFragCombos(t *testing.T) { - c := NewCluster() uri0, err := NewURIFromAddress("host0") if err != nil { t.Fatal(err) @@ -31,8 +33,10 @@ func TestFragCombos(t *testing.T) { if err != nil { t.Fatal(err) } - c.AddNodeBasicSorted(*uri0) - c.AddNodeBasicSorted(*uri1) + + c := NewCluster() + c.addNodeBasicSorted(*uri0) + c.addNodeBasicSorted(*uri1) tests := []struct { idx string @@ -68,3 +72,226 @@ func TestFragCombos(t *testing.T) { } } + +// newIndexWithTempPath returns a new instance of Index. +func newIndexWithTempPath(name string) *Index { + path, err := ioutil.TempDir("", "pilosa-index-") + if err != nil { + panic(err) + } + index, err := NewIndex(path, name) + if err != nil { + panic(err) + } + return index +} + +// Ensure that fragSources creates the correct fragment mapping. +func TestFragSources(t *testing.T) { + + uri0, err := NewURIFromAddress("host0") + if err != nil { + t.Fatal(err) + } + uri1, err := NewURIFromAddress("host1") + if err != nil { + t.Fatal(err) + } + uri2, err := NewURIFromAddress("host2") + if err != nil { + t.Fatal(err) + } + uri3, err := NewURIFromAddress("host3") + if err != nil { + t.Fatal(err) + } + + c1 := NewCluster() + c1.ReplicaN = 1 + c1.addNodeBasicSorted(*uri0) + c1.addNodeBasicSorted(*uri1) + + c2 := NewCluster() + c2.ReplicaN = 1 + c2.addNodeBasicSorted(*uri0) + c2.addNodeBasicSorted(*uri1) + c2.addNodeBasicSorted(*uri2) + + c3 := NewCluster() + c3.ReplicaN = 2 + c3.addNodeBasicSorted(*uri0) + c3.addNodeBasicSorted(*uri1) + + c4 := NewCluster() + c4.ReplicaN = 2 + c4.addNodeBasicSorted(*uri0) + c4.addNodeBasicSorted(*uri1) + c4.addNodeBasicSorted(*uri2) + + c5 := NewCluster() + c5.ReplicaN = 2 + c5.addNodeBasicSorted(*uri0) + c5.addNodeBasicSorted(*uri1) + c5.addNodeBasicSorted(*uri2) + c5.addNodeBasicSorted(*uri3) + + idx := newIndexWithTempPath("i") + frame, err := idx.CreateFrameIfNotExists("f", FrameOptions{}) + if err != nil { + t.Fatal(err) + } + _, err = frame.SetBit("standard", 1, 101, nil) + if err != nil { + t.Fatal(err) + } + _, err = frame.SetBit("standard", 1, 1300000, nil) + if err != nil { + t.Fatal(err) + } + _, err = frame.SetBit("standard", 1, 2600000, nil) + if err != nil { + t.Fatal(err) + } + _, err = frame.SetBit("standard", 1, 3900000, nil) + if err != nil { + t.Fatal(err) + } + + tests := []struct { + from *Cluster + to *Cluster + idx *Index + expected map[URI][]*internal.ResizeSource + err string + }{ + { + from: c1, + to: c2, + idx: idx, + expected: map[URI][]*internal.ResizeSource{ + URI{"http", "host0", 10101}: []*internal.ResizeSource{}, + URI{"http", "host1", 10101}: []*internal.ResizeSource{}, + URI{"http", "host2", 10101}: []*internal.ResizeSource{ + {&internal.URI{"http", "host0", 10101}, "i", "f", "standard", uint64(0)}, + {&internal.URI{"http", "host1", 10101}, "i", "f", "standard", uint64(2)}, + }, + }, + err: "", + }, + { + from: c4, + to: c3, + idx: idx, + expected: map[URI][]*internal.ResizeSource{ + URI{"http", "host0", 10101}: []*internal.ResizeSource{ + {&internal.URI{"http", "host1", 10101}, "i", "f", "standard", uint64(1)}, + }, + URI{"http", "host1", 10101}: []*internal.ResizeSource{ + {&internal.URI{"http", "host0", 10101}, "i", "f", "standard", uint64(0)}, + {&internal.URI{"http", "host0", 10101}, "i", "f", "standard", uint64(2)}, + }, + }, + err: "", + }, + { + from: c5, + to: c4, + idx: idx, + expected: map[URI][]*internal.ResizeSource{ + URI{"http", "host0", 10101}: []*internal.ResizeSource{ + {&internal.URI{"http", "host2", 10101}, "i", "f", "standard", uint64(0)}, + {&internal.URI{"http", "host2", 10101}, "i", "f", "standard", uint64(2)}, + }, + URI{"http", "host1", 10101}: []*internal.ResizeSource{ + {&internal.URI{"http", "host0", 10101}, "i", "f", "standard", uint64(3)}, + }, + URI{"http", "host2", 10101}: []*internal.ResizeSource{}, + }, + err: "", + }, + { + from: c2, + to: c4, + idx: idx, + expected: nil, + err: "clusters are the same size", + }, + { + from: c1, + to: c5, + idx: idx, + expected: nil, + err: "adding more than one node at a time is not supported", + }, + { + from: c5, + to: c1, + idx: idx, + expected: nil, + err: "removing more than one node at a time is not supported", + }, + } + for _, test := range tests { + + actual, err := (test.from).fragSources(test.to, test.idx) + if test.err != "" { + if err.Error() != test.err { + t.Fatalf("expected error: %s", test.err) + } + } else { + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(actual, test.expected) { + t.Errorf("expected: %v, but got: %v", test.expected, actual) + } + } + } +} + +// Ensure that fragSources creates the correct fragment mapping. +func TestResizeJob(t *testing.T) { + + uri0, err := NewURIFromAddress("host0") + if err != nil { + t.Fatal(err) + } + uri1, err := NewURIFromAddress("host1") + if err != nil { + t.Fatal(err) + } + uri2, err := NewURIFromAddress("host2") + if err != nil { + t.Fatal(err) + } + + tests := []struct { + existingURIs []URI + uri URI + action string + expectedURIs map[URI]bool + }{ + { + existingURIs: []URI{*uri0, *uri1}, + uri: *uri2, + action: ResizeJobActionAdd, + expectedURIs: map[URI]bool{*uri0: false, *uri1: false, *uri2: false}, + }, + { + existingURIs: []URI{*uri0, *uri1, *uri2}, + uri: *uri2, + action: ResizeJobActionRemove, + expectedURIs: map[URI]bool{*uri0: false, *uri1: false}, + }, + } + for _, test := range tests { + + actual := NewResizeJob(test.existingURIs, test.uri, test.action) + if err != nil { + t.Fatal(err) + } + if !reflect.DeepEqual(actual.URIs, test.expectedURIs) { + t.Errorf("expected: %v, but got: %v", test.expectedURIs, actual.URIs) + } + } +} diff --git a/cluster_test.go b/cluster_test.go index 0c26b96f2..ca60bf15d 100644 --- a/cluster_test.go +++ b/cluster_test.go @@ -23,7 +23,6 @@ import ( "github.com/davecgh/go-spew/spew" "github.com/pilosa/pilosa" - "github.com/pilosa/pilosa/internal" "github.com/pilosa/pilosa/test" ) @@ -231,80 +230,6 @@ func TestCluster_Topology(t *testing.T) { }) } -// Ensure DataDiff can generate the expected source map. -func TestCluster_Resize(t *testing.T) { - // Given two clusters, ensure DataDiff can determine the sources of data - // needed in order to respond to queries. - t.Run("DataDiff", func(t *testing.T) { - - // Holder - h1 := test.NewHolder() - i, err := h1.CreateIndex("i", pilosa.IndexOptions{}) - if err != nil { - t.Fatal(err) - } - f, err := i.CreateFrame("f", pilosa.FrameOptions{InverseEnabled: true}) - if err != nil { - t.Fatal(err) - } - _, err = f.CreateViewIfNotExists("v") - if err != nil { - t.Fatal(err) - } - _, err = f.CreateViewIfNotExists("inverse") - if err != nil { - t.Fatal(err) - } - - // Set max slices. - i.SetRemoteMaxSlice(2) - i.SetRemoteMaxInverseSlice(5) - - // Cluster 1 - c1 := test.NewCluster(3) - c1.ReplicaN = 2 - - // Cluster 2 - c2 := test.NewCluster(4) - c2.ReplicaN = 2 - - u0 := test.NewURIFromHostPort("host0", 0) - u1 := test.NewURIFromHostPort("host1", 0) - u2 := test.NewURIFromHostPort("host2", 0) - u3 := test.NewURIFromHostPort("host3", 0) - - uri0 := u0.Encode() - uri1 := u1.Encode() - uri2 := u2.Encode() - //uri3 := u3.Encode() - - expected := map[pilosa.URI][]*internal.ResizeSource{ - u0: []*internal.ResizeSource{ - {URI: uri1, Index: "i", Frame: "f", View: "inverse", Slice: 5}, - }, - u1: []*internal.ResizeSource{ - {URI: uri2, Index: "i", Frame: "f", View: "v", Slice: 0}, - {URI: uri2, Index: "i", Frame: "f", View: "inverse", Slice: 0}, - }, - u2: []*internal.ResizeSource{ - {URI: uri0, Index: "i", Frame: "f", View: "inverse", Slice: 3}, - }, - u3: []*internal.ResizeSource{ - {URI: uri0, Index: "i", Frame: "f", View: "v", Slice: 1}, - {URI: uri1, Index: "i", Frame: "f", View: "v", Slice: 2}, - {URI: uri0, Index: "i", Frame: "f", View: "inverse", Slice: 1}, - {URI: uri1, Index: "i", Frame: "f", View: "inverse", Slice: 2}, - {URI: uri1, Index: "i", Frame: "f", View: "inverse", Slice: 5}, - }, - } - - actual := c1.DataDiff(c2, i) - if !reflect.DeepEqual(actual, expected) { - t.Errorf("expected: %v, but got: %v", expected, actual) - } - }) -} - // Ensure that general cluster functionality works as expected. func TestCluster_ResizeStates(t *testing.T) { diff --git a/gossip/gossip.go b/gossip/gossip.go index f1599bf2d..2c5c85faf 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -148,7 +148,7 @@ func NewGossipMemberSet(name string, gossipHost string, gossipPort int, gossipSe g.config.memberlistConfig.BindPort = gossipPort g.config.memberlistConfig.AdvertiseAddr = pilosa.HostToIP(gossipHost) g.config.memberlistConfig.AdvertisePort = gossipPort - g.config.memberlistConfig.PushPullInterval = 0 * time.Second // Default is 15s in DefaultLocalConfig. // TODO travis: change this from 0 + //g.config.memberlistConfig.PushPullInterval = 0 * time.Second // Default is 15s in DefaultLocalConfig. // TODO travis: change this from 0 g.config.memberlistConfig.Delegate = g g.config.memberlistConfig.SecretKey = secretKey g.config.memberlistConfig.Events = server.Cluster.EventReceiver.(memberlist.EventDelegate) diff --git a/handler.go b/handler.go index e0177ba35..943cd13e2 100644 --- a/handler.go +++ b/handler.go @@ -141,6 +141,7 @@ func loadRestricted(router *mux.Router, handler *Handler) { func loadNormal(router *mux.Router, handler *Handler) { router.HandleFunc("/assets/{file}", handler.handleWebUI).Methods("GET") + router.HandleFunc("/cluster/resize/remove-node", handler.handlePostClusterResizeRemoveNode).Methods("POST") router.PathPrefix("/debug/pprof/").Handler(http.DefaultServeMux).Methods("GET") router.HandleFunc("/debug/vars", handler.handleExpvar).Methods("GET") router.HandleFunc("/export", handler.handleGetExport).Methods("GET") @@ -1987,6 +1988,57 @@ type setCoordinatorResponse struct { New *URI `json:"new"` } +// handlePostClusterResizeRemoveNode handles POST /cluster/resize/remove-node request. +func (h *Handler) handlePostClusterResizeRemoveNode(w http.ResponseWriter, r *http.Request) { + // Decode request. + var req removeNodeRequest + err := json.NewDecoder(r.Body).Decode(&req) + if err != nil { + http.Error(w, err.Error(), http.StatusBadRequest) + return + } + + var removeURI *URI + if err := func() error { + removeURI, err = NewURIFromAddress(req.Address) + if err != nil { + return fmt.Errorf("problem with remove node address: %s", err) + } + + // TODO: make sure the address is in the cluster + + // TODO: prevent removing the coordinator node + + // Start the resize process (similar to NodeJoin) + // TODO: this currently blocks. we should leverage listenForJoins() in cluster + // by converted it to a channel of nodeAction {URI, action}. + err := h.Cluster.NodeLeave(*removeURI) + if err != nil { + return err + } + + return nil + }(); err != nil { + http.Error(w, err.Error(), http.StatusInternalServerError) + return + } + + // Encode response. + if err := json.NewEncoder(w).Encode(removeNodeResponse{ + Remove: removeURI, + }); err != nil { + h.logger().Printf("response encoding error: %s", err) + } +} + +type removeNodeRequest struct { + Address string `json:"address"` +} + +type removeNodeResponse struct { + Remove *URI `json:"remove"` +} + // handlePostClusterResizeAbort handles POST /cluster/resize/abort request. func (h *Handler) handlePostClusterResizeAbort(w http.ResponseWriter, r *http.Request) { var msg string diff --git a/index.go b/index.go index 95b09577d..22e6ab062 100644 --- a/index.go +++ b/index.go @@ -780,7 +780,6 @@ func (i *Index) openInputDefinitions() error { return nil } } - } return nil } diff --git a/server.go b/server.go index 55fd62249..c5965912a 100644 --- a/server.go +++ b/server.go @@ -411,6 +411,23 @@ func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error { return nil } + // If this node is still STARTING, don't apply remote status. + // There is an issue where starting up a cluster with existing + // data will error on `flock: resource temporarily unavailable`. + // This is because the ApplySchema creates/opens indexes before + // Holder.Open() has run. When Holder.Open() runs later, the + // fragment files are locked. + // TODO: There is still a race condition where the coordinator + // changes state to NORMAL, broadcasts that to the remote node, + // the remote node receives a `NodeStatus` (with schema) before + // running `Holder.Open()`. In that case, state would be NORMAL, + // meaning this check wouldn't pass, and `Holder.Open()` still + // hasn't run. We may need to track whether `Holder.Open()` has + // run, and use that to determine if we bail here. + if s.Cluster.State == ClusterStateStarting { + return nil + } + // Sync schema. if err := s.Holder.ApplySchema(ns.Schema); err != nil { return err