From a98b862fcae859101068f7b259d2b457361bd5f8 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 20 Nov 2017 18:26:36 -0600 Subject: [PATCH 1/9] add /cluster/resize/remove-node endpoint --- cluster.go | 383 ++++++++++++++++++++++++++++++--------- cluster_internal_test.go | 233 +++++++++++++++++++++++- cluster_test.go | 75 -------- gossip/gossip.go | 2 +- handler.go | 52 ++++++ index.go | 1 - server.go | 17 ++ 7 files changed, 597 insertions(+), 166 deletions(-) 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 From a13db570ccacea63af33444d7b4ff9e462b47eb7 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Tue, 21 Nov 2017 16:23:20 -0600 Subject: [PATCH 2/9] WIP: Wait for nodeState on Holder.Open(). Implement prefect in Cluster so a node can start http listener in a restricted mode. Adjust tests; particularly start test.Holder in state Normal. TODO: - [ ] fix test TestMain_SendReceiveMessage - [ ] add additional tests for `nodeState` --- broadcast.go | 5 + cluster.go | 78 +++++++- cluster_test.go | 20 +- handler.go | 4 +- internal/private.pb.go | 413 ++++++++++++++++++++++++++++++----------- internal/private.proto | 5 + security_manager.go | 8 +- server.go | 52 +++--- test/handler.go | 2 + 9 files changed, 427 insertions(+), 160 deletions(-) diff --git a/broadcast.go b/broadcast.go index 4a5b8e8f7..ec0e53c8d 100644 --- a/broadcast.go +++ b/broadcast.go @@ -126,6 +126,7 @@ const ( MessageTypeResizeInstruction = 10 MessageTypeResizeInstructionComplete = 11 MessageTypeSetCoordinator = 12 + MessageTypeNodeState = 13 ) // MarshalMessage encodes the protobuf message into a byte slice. @@ -156,6 +157,8 @@ func MarshalMessage(m proto.Message) ([]byte, error) { typ = MessageTypeResizeInstructionComplete case *internal.SetCoordinatorMessage: typ = MessageTypeSetCoordinator + case *internal.NodeStateMessage: + typ = MessageTypeNodeState default: return nil, fmt.Errorf("message type not implemented for marshalling: %s", reflect.TypeOf(obj)) } @@ -196,6 +199,8 @@ func UnmarshalMessage(buf []byte) (proto.Message, error) { m = &internal.ResizeInstructionComplete{} case MessageTypeSetCoordinator: m = &internal.SetCoordinatorMessage{} + case MessageTypeNodeState: + m = &internal.NodeStateMessage{} default: return nil, fmt.Errorf("invalid message type: %d", typ) } diff --git a/cluster.go b/cluster.go index befc66c1a..171615b4d 100644 --- a/cluster.go +++ b/cluster.go @@ -46,6 +46,10 @@ const ( ClusterStateNormal = "NORMAL" ClusterStateResizing = "RESIZING" + // NodeState represents the state of a node during startup. + NodeStateLoading = "LOADING" + NodeStateReady = "READY" + // ResizeJob states. ResizeJobStateRunning = "RUNNING" // Final states. @@ -285,6 +289,41 @@ func (c *Cluster) setState(state string) { c.State = state } +func (c *Cluster) setNodeState(state string) { + if c.IsCoordinator() { + c.Topology.nodeStates[c.URI] = state + return + } + + // Send node state to coordinator. + ns := &internal.NodeStateMessage{ + URI: c.URI.Encode(), + State: state, + } + + node := &Node{ + URI: c.Coordinator, + } + if err := c.Broadcaster.SendTo(node, ns); err != nil { + c.logger().Printf("sending node state error: err=%s", err) + } +} + +func (c *Cluster) ReceiveNodeState(uri URI, state string) error { + if !c.IsCoordinator() { + return nil + } + + c.Topology.nodeStates[uri] = state + + // Set cluster state to NORMAL. + if c.haveTopologyAgreement() && c.allNodesReady() { + return c.setStateAndBroadcast(ClusterStateNormal) + } + + return nil +} + // localNode is not being used. //func (c *Cluster) localNode() *Node { // return c.NodeByURI(c.URI) @@ -693,6 +732,15 @@ func (c *Cluster) haveTopologyAgreement() bool { return URISlicesAreEqual(c.Topology.NodeSet, c.NodeSet()) } +func (c *Cluster) allNodesReady() bool { + for _, uri := range c.Topology.NodeSet { + if c.Topology.nodeStates[uri] != NodeStateReady { + return false + } + } + return true +} + func (c *Cluster) handleNodeAction(nodeAction nodeAction) error { j, err := c.generateResizeJob(nodeAction) if err != nil { @@ -1130,10 +1178,16 @@ func (u NodeSet) ToStrings() []string { type Topology struct { mu sync.RWMutex NodeSet []URI + + // nodeStates holds the state of each node according to + // the coordinator. Used during startup and data load. + nodeStates map[URI]string } func NewTopology() *Topology { - return &Topology{} + return &Topology{ + nodeStates: make(map[URI]string), + } } // ContainsURI returns true if uri matches one of the topology's uris. @@ -1241,9 +1295,9 @@ func decodeTopology(topology *internal.Topology) (*Topology, error) { return nil, nil } - t := &Topology{ - NodeSet: decodeURIs(topology.NodeSet), - } + t := NewTopology() + t.NodeSet = decodeURIs(topology.NodeSet) + return t, nil } @@ -1302,9 +1356,19 @@ func (c *Cluster) nodeJoin(uri URI) error { 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() { + // Only change to normal if there is no existing data. Otherwise, + // the coordinator needs to wait to receive READY messages (nodeStates) + // from remote nodes before setting the cluster to state NORMAL. + if !c.Holder.HasData() { + // 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 + } + + if c.haveTopologyAgreement() && c.allNodesReady() { return c.setStateAndBroadcast(ClusterStateNormal) } diff --git a/cluster_test.go b/cluster_test.go index ca60bf15d..b13638bf4 100644 --- a/cluster_test.go +++ b/cluster_test.go @@ -253,8 +253,8 @@ func TestCluster_ResizeStates(t *testing.T) { } // Verify topology file. - if !reflect.DeepEqual(node.Topology, expectedTop) { - t.Errorf("expected topology: %v, but got: %v", expectedTop, node.Topology) + if !reflect.DeepEqual(node.Topology.NodeSet, expectedTop.NodeSet) { + t.Errorf("expected topology: %v, but got: %v", expectedTop.NodeSet, node.Topology.NodeSet) } // Close TestCluster. @@ -344,10 +344,10 @@ func TestCluster_ResizeStates(t *testing.T) { } // 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) + if !reflect.DeepEqual(node0.Topology.NodeSet, expectedTop.NodeSet) { + t.Errorf("expected node0 topology: %v, but got: %v", expectedTop.NodeSet, node0.Topology.NodeSet) + } else if !reflect.DeepEqual(node1.Topology.NodeSet, expectedTop.NodeSet) { + t.Errorf("expected node1 topology: %v, but got: %v", expectedTop.NodeSet, node1.Topology.NodeSet) } // Close TestCluster. @@ -436,10 +436,10 @@ func TestCluster_ResizeStates(t *testing.T) { } // 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) + if !reflect.DeepEqual(node0.Topology.NodeSet, expectedTop.NodeSet) { + t.Errorf("expected node0 topology: %v, but got: %v", expectedTop.NodeSet, node0.Topology.NodeSet) + } else if !reflect.DeepEqual(node1.Topology.NodeSet, expectedTop.NodeSet) { + t.Errorf("expected node1 topology: %v, but got: %v", expectedTop.NodeSet, node1.Topology.NodeSet) } // Verify that node-1 contains the fragment (i/f/standard/1) transferred from node-0. diff --git a/handler.go b/handler.go index 943cd13e2..1ef7be7f9 100644 --- a/handler.go +++ b/handler.go @@ -109,7 +109,7 @@ func BuildRouters(handler *Handler) { loadRestricted(router, handler) handler.RestrictedRouter = router - handler.SetNormal() + handler.SetRestricted() } // SetNormal is a method of the SecurityManager interface which provides normal URI routing. @@ -184,7 +184,7 @@ func loadNormal(router *mux.Router, handler *Handler) { } func (h *Handler) reportRestricted(w http.ResponseWriter, r *http.Request) { - http.Error(w, "not allowed during resize", http.StatusMethodNotAllowed) + http.Error(w, fmt.Sprintf("not allowed in cluster state %s", h.Cluster.State), http.StatusMethodNotAllowed) } func (h *Handler) methodNotAllowedHandler(w http.ResponseWriter, r *http.Request) { diff --git a/internal/private.pb.go b/internal/private.pb.go index 172bc0be2..0b43b5cc9 100644 --- a/internal/private.pb.go +++ b/internal/private.pb.go @@ -29,6 +29,7 @@ CreateInputDefinitionMessage DeleteInputDefinitionMessage URI + NodeStateMessage NodeStatus ClusterStatus Field @@ -670,6 +671,30 @@ func (m *URI) GetPort() uint32 { return 0 } +type NodeStateMessage struct { + URI *URI `protobuf:"bytes,1,opt,name=URI" json:"URI,omitempty"` + State string `protobuf:"bytes,2,opt,name=State,proto3" json:"State,omitempty"` +} + +func (m *NodeStateMessage) Reset() { *m = NodeStateMessage{} } +func (m *NodeStateMessage) String() string { return proto.CompactTextString(m) } +func (*NodeStateMessage) ProtoMessage() {} +func (*NodeStateMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{21} } + +func (m *NodeStateMessage) GetURI() *URI { + if m != nil { + return m.URI + } + return nil +} + +func (m *NodeStateMessage) GetState() string { + if m != nil { + return m.State + } + return "" +} + type NodeStatus struct { URI *URI `protobuf:"bytes,1,opt,name=URI" json:"URI,omitempty"` MaxSlices *MaxSlices `protobuf:"bytes,2,opt,name=MaxSlices" json:"MaxSlices,omitempty"` @@ -679,7 +704,7 @@ type NodeStatus struct { func (m *NodeStatus) Reset() { *m = NodeStatus{} } func (m *NodeStatus) String() string { return proto.CompactTextString(m) } func (*NodeStatus) ProtoMessage() {} -func (*NodeStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{21} } +func (*NodeStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{22} } func (m *NodeStatus) GetURI() *URI { if m != nil { @@ -710,7 +735,7 @@ type ClusterStatus struct { func (m *ClusterStatus) Reset() { *m = ClusterStatus{} } func (m *ClusterStatus) String() string { return proto.CompactTextString(m) } func (*ClusterStatus) ProtoMessage() {} -func (*ClusterStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{22} } +func (*ClusterStatus) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{23} } func (m *ClusterStatus) GetState() string { if m != nil { @@ -736,7 +761,7 @@ type Field struct { func (m *Field) Reset() { *m = Field{} } func (m *Field) String() string { return proto.CompactTextString(m) } func (*Field) ProtoMessage() {} -func (*Field) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{23} } +func (*Field) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{24} } func (m *Field) GetName() string { if m != nil { @@ -775,7 +800,7 @@ type DeleteViewMessage struct { func (m *DeleteViewMessage) Reset() { *m = DeleteViewMessage{} } func (m *DeleteViewMessage) String() string { return proto.CompactTextString(m) } func (*DeleteViewMessage) ProtoMessage() {} -func (*DeleteViewMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{24} } +func (*DeleteViewMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{25} } func (m *DeleteViewMessage) GetIndex() string { if m != nil { @@ -809,7 +834,7 @@ type ResizeInstruction struct { func (m *ResizeInstruction) Reset() { *m = ResizeInstruction{} } func (m *ResizeInstruction) String() string { return proto.CompactTextString(m) } func (*ResizeInstruction) ProtoMessage() {} -func (*ResizeInstruction) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{25} } +func (*ResizeInstruction) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{26} } func (m *ResizeInstruction) GetJobID() int64 { if m != nil { @@ -857,7 +882,7 @@ type ResizeSource struct { func (m *ResizeSource) Reset() { *m = ResizeSource{} } func (m *ResizeSource) String() string { return proto.CompactTextString(m) } func (*ResizeSource) ProtoMessage() {} -func (*ResizeSource) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{26} } +func (*ResizeSource) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{27} } func (m *ResizeSource) GetURI() *URI { if m != nil { @@ -904,7 +929,7 @@ func (m *ResizeInstructionComplete) Reset() { *m = ResizeInstructionComp func (m *ResizeInstructionComplete) String() string { return proto.CompactTextString(m) } func (*ResizeInstructionComplete) ProtoMessage() {} func (*ResizeInstructionComplete) Descriptor() ([]byte, []int) { - return fileDescriptorPrivate, []int{27} + return fileDescriptorPrivate, []int{28} } func (m *ResizeInstructionComplete) GetJobID() int64 { @@ -936,7 +961,7 @@ type SetCoordinatorMessage struct { func (m *SetCoordinatorMessage) Reset() { *m = SetCoordinatorMessage{} } func (m *SetCoordinatorMessage) String() string { return proto.CompactTextString(m) } func (*SetCoordinatorMessage) ProtoMessage() {} -func (*SetCoordinatorMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{28} } +func (*SetCoordinatorMessage) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{29} } func (m *SetCoordinatorMessage) GetOld() *URI { if m != nil { @@ -959,7 +984,7 @@ type Topology struct { func (m *Topology) Reset() { *m = Topology{} } func (m *Topology) String() string { return proto.CompactTextString(m) } func (*Topology) ProtoMessage() {} -func (*Topology) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{29} } +func (*Topology) Descriptor() ([]byte, []int) { return fileDescriptorPrivate, []int{30} } func (m *Topology) GetNodeSet() []*URI { if m != nil { @@ -990,6 +1015,7 @@ func init() { proto.RegisterType((*CreateInputDefinitionMessage)(nil), "internal.CreateInputDefinitionMessage") proto.RegisterType((*DeleteInputDefinitionMessage)(nil), "internal.DeleteInputDefinitionMessage") proto.RegisterType((*URI)(nil), "internal.URI") + proto.RegisterType((*NodeStateMessage)(nil), "internal.NodeStateMessage") proto.RegisterType((*NodeStatus)(nil), "internal.NodeStatus") proto.RegisterType((*ClusterStatus)(nil), "internal.ClusterStatus") proto.RegisterType((*Field)(nil), "internal.Field") @@ -1848,6 +1874,40 @@ func (m *URI) MarshalTo(dAtA []byte) (int, error) { return i, nil } +func (m *NodeStateMessage) Marshal() (dAtA []byte, err error) { + size := m.Size() + dAtA = make([]byte, size) + n, err := m.MarshalTo(dAtA) + if err != nil { + return nil, err + } + return dAtA[:n], nil +} + +func (m *NodeStateMessage) MarshalTo(dAtA []byte) (int, error) { + var i int + _ = i + var l int + _ = l + if m.URI != nil { + dAtA[i] = 0xa + i++ + i = encodeVarintPrivate(dAtA, i, uint64(m.URI.Size())) + n11, err := m.URI.MarshalTo(dAtA[i:]) + if err != nil { + return 0, err + } + i += n11 + } + if len(m.State) > 0 { + dAtA[i] = 0x12 + i++ + i = encodeVarintPrivate(dAtA, i, uint64(len(m.State))) + i += copy(dAtA[i:], m.State) + } + return i, nil +} + func (m *NodeStatus) Marshal() (dAtA []byte, err error) { size := m.Size() dAtA = make([]byte, size) @@ -1867,32 +1927,32 @@ func (m *NodeStatus) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0xa i++ i = encodeVarintPrivate(dAtA, i, uint64(m.URI.Size())) - n11, err := m.URI.MarshalTo(dAtA[i:]) - if err != nil { - return 0, err - } - i += n11 - } - if m.MaxSlices != nil { - dAtA[i] = 0x12 - i++ - i = encodeVarintPrivate(dAtA, i, uint64(m.MaxSlices.Size())) - n12, err := m.MaxSlices.MarshalTo(dAtA[i:]) + n12, err := m.URI.MarshalTo(dAtA[i:]) if err != nil { return 0, err } i += n12 } - if m.Schema != nil { - dAtA[i] = 0x1a + if m.MaxSlices != nil { + dAtA[i] = 0x12 i++ - i = encodeVarintPrivate(dAtA, i, uint64(m.Schema.Size())) - n13, err := m.Schema.MarshalTo(dAtA[i:]) + i = encodeVarintPrivate(dAtA, i, uint64(m.MaxSlices.Size())) + n13, err := m.MaxSlices.MarshalTo(dAtA[i:]) if err != nil { return 0, err } i += n13 } + if m.Schema != nil { + dAtA[i] = 0x1a + i++ + i = encodeVarintPrivate(dAtA, i, uint64(m.Schema.Size())) + n14, err := m.Schema.MarshalTo(dAtA[i:]) + if err != nil { + return 0, err + } + i += n14 + } return i, nil } @@ -2032,21 +2092,21 @@ func (m *ResizeInstruction) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0x12 i++ i = encodeVarintPrivate(dAtA, i, uint64(m.URI.Size())) - n14, err := m.URI.MarshalTo(dAtA[i:]) + n15, err := m.URI.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n14 + i += n15 } if m.Coordinator != nil { dAtA[i] = 0x1a i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Coordinator.Size())) - n15, err := m.Coordinator.MarshalTo(dAtA[i:]) + n16, err := m.Coordinator.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n15 + i += n16 } if len(m.Sources) > 0 { for _, msg := range m.Sources { @@ -2064,11 +2124,11 @@ func (m *ResizeInstruction) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0x2a i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Schema.Size())) - n16, err := m.Schema.MarshalTo(dAtA[i:]) + n17, err := m.Schema.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n16 + i += n17 } return i, nil } @@ -2092,11 +2152,11 @@ func (m *ResizeSource) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0xa i++ i = encodeVarintPrivate(dAtA, i, uint64(m.URI.Size())) - n17, err := m.URI.MarshalTo(dAtA[i:]) + n18, err := m.URI.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n17 + i += n18 } if len(m.Index) > 0 { dAtA[i] = 0x12 @@ -2148,11 +2208,11 @@ func (m *ResizeInstructionComplete) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0x12 i++ i = encodeVarintPrivate(dAtA, i, uint64(m.URI.Size())) - n18, err := m.URI.MarshalTo(dAtA[i:]) + n19, err := m.URI.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n18 + i += n19 } if len(m.Error) > 0 { dAtA[i] = 0x1a @@ -2182,21 +2242,21 @@ func (m *SetCoordinatorMessage) MarshalTo(dAtA []byte) (int, error) { dAtA[i] = 0xa i++ i = encodeVarintPrivate(dAtA, i, uint64(m.Old.Size())) - n19, err := m.Old.MarshalTo(dAtA[i:]) + n20, err := m.Old.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n19 + i += n20 } if m.New != nil { dAtA[i] = 0x12 i++ i = encodeVarintPrivate(dAtA, i, uint64(m.New.Size())) - n20, err := m.New.MarshalTo(dAtA[i:]) + n21, err := m.New.MarshalTo(dAtA[i:]) if err != nil { return 0, err } - i += n20 + i += n21 } return i, nil } @@ -2613,6 +2673,20 @@ func (m *URI) Size() (n int) { return n } +func (m *NodeStateMessage) Size() (n int) { + var l int + _ = l + if m.URI != nil { + l = m.URI.Size() + n += 1 + l + sovPrivate(uint64(l)) + } + l = len(m.State) + if l > 0 { + n += 1 + l + sovPrivate(uint64(l)) + } + return n +} + func (m *NodeStatus) Size() (n int) { var l int _ = l @@ -5705,6 +5779,118 @@ func (m *URI) Unmarshal(dAtA []byte) error { } return nil } +func (m *NodeStateMessage) Unmarshal(dAtA []byte) error { + l := len(dAtA) + iNdEx := 0 + for iNdEx < l { + preIndex := iNdEx + var wire uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + wire |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + fieldNum := int32(wire >> 3) + wireType := int(wire & 0x7) + if wireType == 4 { + return fmt.Errorf("proto: NodeStateMessage: wiretype end group for non-group") + } + if fieldNum <= 0 { + return fmt.Errorf("proto: NodeStateMessage: illegal tag %d (wire type %d)", fieldNum, wire) + } + switch fieldNum { + case 1: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field URI", wireType) + } + var msglen int + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + msglen |= (int(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + if msglen < 0 { + return ErrInvalidLengthPrivate + } + postIndex := iNdEx + msglen + if postIndex > l { + return io.ErrUnexpectedEOF + } + if m.URI == nil { + m.URI = &URI{} + } + if err := m.URI.Unmarshal(dAtA[iNdEx:postIndex]); err != nil { + return err + } + iNdEx = postIndex + case 2: + if wireType != 2 { + return fmt.Errorf("proto: wrong wireType = %d for field State", wireType) + } + var stringLen uint64 + for shift := uint(0); ; shift += 7 { + if shift >= 64 { + return ErrIntOverflowPrivate + } + if iNdEx >= l { + return io.ErrUnexpectedEOF + } + b := dAtA[iNdEx] + iNdEx++ + stringLen |= (uint64(b) & 0x7F) << shift + if b < 0x80 { + break + } + } + intStringLen := int(stringLen) + if intStringLen < 0 { + return ErrInvalidLengthPrivate + } + postIndex := iNdEx + intStringLen + if postIndex > l { + return io.ErrUnexpectedEOF + } + m.State = string(dAtA[iNdEx:postIndex]) + iNdEx = postIndex + default: + iNdEx = preIndex + skippy, err := skipPrivate(dAtA[iNdEx:]) + if err != nil { + return err + } + if skippy < 0 { + return ErrInvalidLengthPrivate + } + if (iNdEx + skippy) > l { + return io.ErrUnexpectedEOF + } + iNdEx += skippy + } + } + + if iNdEx > l { + return io.ErrUnexpectedEOF + } + return nil +} func (m *NodeStatus) Unmarshal(dAtA []byte) error { l := len(dAtA) iNdEx := 0 @@ -7071,79 +7257,80 @@ var ( func init() { proto.RegisterFile("private.proto", fileDescriptorPrivate) } var fileDescriptorPrivate = []byte{ - // 1179 bytes of a gzipped FileDescriptorProto - 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x9c, 0x57, 0x4f, 0x6f, 0xe3, 0x44, - 0x14, 0xc7, 0x71, 0x92, 0x26, 0x2f, 0x9b, 0xdd, 0xec, 0xd0, 0xad, 0xd2, 0xaa, 0x4a, 0xc3, 0x1c, - 0x68, 0x59, 0x89, 0x02, 0xad, 0x84, 0xa0, 0x08, 0x09, 0xb6, 0xc9, 0x6a, 0x0d, 0xb4, 0x5d, 0x26, - 0xdd, 0x45, 0x70, 0x40, 0x9a, 0x26, 0x43, 0x6b, 0xd5, 0xb1, 0x83, 0x3d, 0x69, 0x9b, 0x3d, 0x70, - 0x83, 0x03, 0x7c, 0x01, 0xee, 0x7c, 0x19, 0x8e, 0xdc, 0xb8, 0xa2, 0xf2, 0x21, 0x90, 0xb8, 0x80, - 0xe6, 0x79, 0xc6, 0x76, 0xfe, 0x35, 0x6a, 0x6f, 0x7e, 0x6f, 0xde, 0x7b, 0xf3, 0x9b, 0xdf, 0xfb, - 0x33, 0x63, 0xa8, 0x0e, 0x42, 0xf7, 0x82, 0x4b, 0xb1, 0x3d, 0x08, 0x03, 0x19, 0x90, 0x92, 0xeb, - 0x4b, 0x11, 0xfa, 0xdc, 0xa3, 0x47, 0x50, 0x76, 0xfc, 0x9e, 0xb8, 0x3a, 0x10, 0x92, 0x93, 0x26, - 0x54, 0xf6, 0x03, 0x6f, 0xd8, 0xf7, 0xbf, 0xe0, 0x27, 0xc2, 0xab, 0x5b, 0x4d, 0x6b, 0xab, 0xcc, - 0xb2, 0x2a, 0x65, 0x71, 0xec, 0xf6, 0xc5, 0x97, 0x43, 0xee, 0xcb, 0x61, 0xbf, 0x9e, 0x8b, 0x2d, - 0x32, 0x2a, 0xfa, 0xaf, 0x05, 0xe5, 0xa7, 0x21, 0xef, 0x0b, 0x8c, 0xb8, 0x06, 0x25, 0x16, 0x5c, - 0x66, 0xc3, 0x25, 0x32, 0x79, 0x13, 0xee, 0x3b, 0xfe, 0x85, 0x08, 0x23, 0xd1, 0xf6, 0xf9, 0x89, - 0x27, 0x7a, 0x18, 0xae, 0xc4, 0x26, 0xb4, 0x64, 0x1d, 0xca, 0xfb, 0xbc, 0x7b, 0x26, 0x8e, 0x47, - 0x03, 0x51, 0xb7, 0x31, 0x48, 0xaa, 0x48, 0x56, 0x3b, 0xee, 0x2b, 0x51, 0xcf, 0x37, 0xad, 0xad, - 0x2a, 0x4b, 0x15, 0x93, 0x78, 0x0b, 0x53, 0x78, 0x09, 0x85, 0x7b, 0x8c, 0xfb, 0xa7, 0x09, 0x86, - 0x22, 0x62, 0x18, 0xd3, 0x91, 0x4d, 0x28, 0x3e, 0x75, 0x85, 0xd7, 0x8b, 0xea, 0x4b, 0x4d, 0x7b, - 0xab, 0xb2, 0xf3, 0x60, 0xdb, 0xf0, 0xb7, 0x8d, 0x7a, 0xa6, 0x97, 0x29, 0x85, 0xfb, 0x4e, 0x7f, - 0x10, 0x84, 0x92, 0x89, 0x68, 0x10, 0xf8, 0x91, 0x20, 0x35, 0xb0, 0xdb, 0x61, 0xa8, 0xcf, 0xae, - 0x3e, 0xe9, 0x0f, 0x50, 0x7b, 0xe2, 0x05, 0xdd, 0xf3, 0x16, 0x97, 0x9c, 0x89, 0xef, 0x87, 0x22, - 0x92, 0x64, 0x19, 0x0a, 0x98, 0x05, 0x6d, 0x17, 0x0b, 0x4a, 0x8b, 0x4c, 0x6a, 0x9a, 0x63, 0x41, - 0x69, 0xd1, 0x1f, 0xa9, 0xc8, 0xb3, 0x58, 0x50, 0xda, 0x8e, 0xe7, 0x76, 0x63, 0x0a, 0xf2, 0x2c, - 0x16, 0x08, 0x81, 0xfc, 0x4b, 0x57, 0x5c, 0xea, 0x73, 0xe3, 0x37, 0x75, 0xe0, 0x61, 0x66, 0x7f, - 0x0d, 0x73, 0x05, 0x8a, 0x2c, 0xb8, 0x74, 0x5a, 0x51, 0xdd, 0x6a, 0xda, 0x5b, 0x79, 0xa6, 0x25, - 0x64, 0x17, 0xd3, 0xaf, 0x96, 0x72, 0xb8, 0x94, 0x2a, 0xe8, 0x2a, 0x14, 0x90, 0x6a, 0x75, 0xca, - 0xd4, 0x57, 0x7d, 0xd2, 0xff, 0x2c, 0x28, 0x1f, 0xf0, 0x2b, 0x84, 0x11, 0x91, 0x8f, 0xa1, 0xd4, - 0x91, 0xdc, 0xef, 0xf1, 0xb0, 0x87, 0x46, 0x95, 0x9d, 0x37, 0x52, 0x0a, 0x13, 0xb3, 0x6d, 0x63, - 0xd3, 0xf6, 0x65, 0x38, 0x62, 0x89, 0x0b, 0xd9, 0x83, 0x25, 0x5d, 0x13, 0x88, 0xa1, 0xb2, 0xd3, - 0x9c, 0xe5, 0x9d, 0x94, 0x8d, 0x72, 0x36, 0x0e, 0x6b, 0x1f, 0x41, 0x75, 0x2c, 0xac, 0xc2, 0x7a, - 0x2e, 0x46, 0x26, 0x23, 0xe7, 0x62, 0xa4, 0xb8, 0xbb, 0xe0, 0xde, 0x30, 0xe6, 0x39, 0xcf, 0x62, - 0x61, 0x2f, 0xf7, 0x81, 0xb5, 0xb6, 0x07, 0xf7, 0xb2, 0x51, 0x6f, 0xe3, 0x4b, 0xbf, 0x05, 0xb2, - 0x1f, 0x0a, 0x2e, 0x05, 0xc2, 0x3b, 0x10, 0x51, 0xc4, 0x4f, 0xc5, 0xfc, 0x4c, 0xc7, 0xd9, 0xcb, - 0x65, 0xb3, 0xb7, 0x0e, 0x65, 0x27, 0x32, 0x07, 0xb7, 0xb1, 0x2e, 0x53, 0x05, 0x7d, 0x0c, 0xa4, - 0x25, 0x3c, 0x21, 0x85, 0xee, 0xdf, 0x1b, 0xe2, 0xd3, 0x8e, 0xc1, 0xb2, 0xd8, 0x96, 0x6c, 0x42, - 0x5e, 0xb5, 0x2e, 0x42, 0xa9, 0xec, 0xbc, 0x9e, 0x32, 0x9d, 0xcc, 0x09, 0x86, 0x06, 0xd4, 0x35, - 0x41, 0x75, 0xbb, 0x2f, 0x38, 0xe0, 0x8c, 0x52, 0x36, 0x5b, 0xd9, 0x93, 0x5b, 0x25, 0x03, 0x44, - 0x6f, 0xf5, 0x89, 0x39, 0xeb, 0x5d, 0xb7, 0xa2, 0xdf, 0x68, 0xad, 0x6a, 0x89, 0x43, 0xb5, 0x1a, - 0xfb, 0xe0, 0xf7, 0xfc, 0x23, 0x4f, 0xe0, 0x50, 0xb1, 0x55, 0x0f, 0x45, 0x75, 0xbb, 0x69, 0xab, - 0xd8, 0x28, 0xd0, 0x5d, 0x28, 0x76, 0xba, 0x67, 0xa2, 0xcf, 0xc9, 0x5b, 0xaa, 0x50, 0x7b, 0xe2, - 0x4a, 0x44, 0xba, 0xcc, 0x1f, 0x4c, 0xd0, 0xc7, 0xcc, 0x3a, 0xfd, 0xc5, 0xd2, 0xe8, 0xe7, 0x20, - 0x2a, 0xe2, 0xde, 0x51, 0x3d, 0x3f, 0x35, 0x71, 0x94, 0x9e, 0xe9, 0x65, 0xd2, 0x86, 0x9a, 0xe3, - 0x0f, 0x86, 0xb2, 0x25, 0xbe, 0x73, 0x7d, 0x57, 0xba, 0x81, 0x1f, 0xd5, 0x8b, 0xe8, 0xb2, 0x9a, - 0xdd, 0x7a, 0xcc, 0x82, 0x4d, 0xb9, 0xd0, 0x9f, 0x2c, 0x78, 0x30, 0xa1, 0x5c, 0x80, 0x2b, 0x77, - 0x33, 0xae, 0xf7, 0x93, 0x91, 0x69, 0xa3, 0x61, 0x63, 0x2e, 0x9a, 0xf1, 0x09, 0xfa, 0x9b, 0x05, - 0xcb, 0xb3, 0x0c, 0x66, 0xa2, 0x69, 0x00, 0x3c, 0x0f, 0xdd, 0x3e, 0x0f, 0x47, 0x9f, 0x8b, 0x91, - 0xbe, 0x3d, 0x32, 0x1a, 0xf2, 0x15, 0xac, 0x4c, 0xc4, 0xfa, 0xb4, 0x1b, 0x53, 0x14, 0x83, 0xda, - 0x98, 0x0b, 0x2a, 0xb6, 0x63, 0x73, 0xdc, 0xe9, 0x3f, 0x16, 0x3c, 0x9a, 0xb9, 0x94, 0x56, 0x9f, - 0x95, 0x2d, 0xf4, 0xc7, 0x50, 0x7b, 0xa9, 0x06, 0x43, 0x4b, 0x44, 0xd2, 0xf5, 0xb9, 0xb2, 0xd4, - 0xe5, 0x39, 0xa5, 0x27, 0x0e, 0x94, 0x50, 0x77, 0xc0, 0x07, 0x1a, 0xe6, 0xdb, 0x0b, 0x60, 0x6e, - 0x1b, 0x7b, 0x3d, 0x37, 0x8d, 0xa8, 0xc0, 0xe0, 0x1c, 0x37, 0x97, 0x02, 0x0a, 0x6a, 0x22, 0x8e, - 0x39, 0xdc, 0x6a, 0xaa, 0x05, 0xb0, 0x6e, 0x26, 0xc9, 0x18, 0x92, 0x9b, 0x7b, 0xf2, 0x43, 0x80, - 0xd4, 0x54, 0xb7, 0xfb, 0x0d, 0xf5, 0x99, 0x31, 0xa6, 0xcf, 0x60, 0xdd, 0x8c, 0xb9, 0x5b, 0x6c, - 0x68, 0xaa, 0x25, 0x97, 0x56, 0x0b, 0x6d, 0x83, 0xfd, 0x82, 0x39, 0xea, 0xaa, 0xc3, 0x6e, 0x35, - 0x29, 0xd2, 0x92, 0x72, 0x79, 0x16, 0x44, 0xd2, 0xb8, 0xa8, 0x6f, 0xa5, 0x7b, 0x1e, 0x84, 0x12, - 0x11, 0x57, 0x19, 0x7e, 0xd3, 0x9f, 0x2d, 0x80, 0xc3, 0xa0, 0x27, 0x3a, 0x92, 0xcb, 0x61, 0x44, - 0x36, 0x30, 0x2a, 0xc6, 0xaa, 0xec, 0x54, 0xd3, 0x33, 0xbd, 0x60, 0x0e, 0xc3, 0xfd, 0xde, 0xcb, - 0x5c, 0x84, 0xd3, 0x13, 0x26, 0x59, 0x62, 0x99, 0xeb, 0x72, 0xcb, 0x0c, 0x14, 0x4d, 0x55, 0x2d, - 0xb5, 0x8f, 0xf5, 0x1a, 0x34, 0xa7, 0x87, 0x50, 0xdd, 0xf7, 0x86, 0x91, 0x14, 0xa1, 0x86, 0xa3, - 0x6e, 0x12, 0xc9, 0x65, 0x52, 0x7f, 0x28, 0x90, 0x4d, 0x58, 0x42, 0xc8, 0x42, 0xea, 0xbe, 0x9d, - 0x00, 0x6a, 0x56, 0x69, 0x07, 0x0a, 0xf3, 0xdb, 0x8d, 0x40, 0x1e, 0xdf, 0x60, 0x9a, 0x21, 0x7c, - 0x7e, 0xd5, 0xc0, 0x3e, 0x70, 0xe3, 0x94, 0xda, 0x4c, 0x7d, 0xa2, 0x86, 0x5f, 0x61, 0xc9, 0x29, - 0x0d, 0x57, 0xb7, 0xcf, 0xc3, 0x38, 0x85, 0x6a, 0x5c, 0xde, 0xe5, 0x9e, 0x30, 0xcf, 0x18, 0x3b, - 0xf3, 0x8c, 0xf9, 0xd3, 0x82, 0x87, 0x4c, 0x44, 0xee, 0x2b, 0xe1, 0xf8, 0x91, 0x0c, 0x87, 0x49, - 0xfb, 0x7d, 0x16, 0x9c, 0x38, 0x2d, 0x8c, 0x6a, 0xb3, 0x58, 0x30, 0x39, 0xca, 0xcd, 0xcd, 0xd1, - 0x3b, 0xea, 0xe1, 0x1b, 0x84, 0x3d, 0xd5, 0x83, 0x41, 0xa8, 0x59, 0x9f, 0x30, 0xcc, 0x5a, 0x90, - 0x77, 0x61, 0xa9, 0x13, 0x0c, 0xc3, 0x6e, 0x32, 0xa0, 0x57, 0x52, 0xe3, 0x18, 0x55, 0xbc, 0xcc, - 0x8c, 0x59, 0x26, 0xa7, 0x85, 0x05, 0x39, 0xfd, 0xd1, 0x82, 0x7b, 0xd9, 0x18, 0x8b, 0x4b, 0x2c, - 0xe1, 0x32, 0x37, 0x93, 0x4b, 0x7b, 0x16, 0x97, 0xf9, 0x94, 0xcb, 0xf4, 0xf9, 0x51, 0xc8, 0x3c, - 0x3f, 0xe8, 0x19, 0xac, 0x4e, 0x11, 0xbc, 0x1f, 0xf4, 0x07, 0x2a, 0x93, 0x77, 0x25, 0x7a, 0x19, - 0x0a, 0xed, 0x30, 0xd4, 0x14, 0x97, 0x59, 0x2c, 0xd0, 0xaf, 0xe1, 0x51, 0x47, 0xc8, 0x0c, 0xbf, - 0xa6, 0x48, 0x36, 0xc0, 0x3e, 0xf2, 0x7a, 0x73, 0x4e, 0x7e, 0xe4, 0xf5, 0x94, 0xc1, 0xa1, 0xb8, - 0x9c, 0xb3, 0xe1, 0xa1, 0xb8, 0xa4, 0xbb, 0x50, 0x3a, 0x0e, 0x06, 0x81, 0x17, 0x9c, 0x8e, 0xb2, - 0x5d, 0x60, 0xdd, 0xd4, 0x05, 0x4f, 0x6a, 0xbf, 0x5f, 0x37, 0xac, 0x3f, 0xae, 0x1b, 0xd6, 0x5f, - 0xd7, 0x0d, 0xeb, 0xd7, 0xbf, 0x1b, 0xaf, 0x9d, 0x14, 0xf1, 0xbf, 0x69, 0xf7, 0xff, 0x00, 0x00, - 0x00, 0xff, 0xff, 0x80, 0x0a, 0xbc, 0x80, 0x48, 0x0d, 0x00, 0x00, + // 1192 bytes of a gzipped FileDescriptorProto + 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x9c, 0x57, 0x4d, 0x6f, 0x23, 0x45, + 0x13, 0x7e, 0xc7, 0x63, 0x3b, 0x76, 0x79, 0xbd, 0xeb, 0xf4, 0x9b, 0x8d, 0x9c, 0x28, 0xf2, 0x9a, + 0x3e, 0x90, 0xb0, 0x12, 0x01, 0x12, 0x09, 0x41, 0x10, 0x12, 0x6c, 0xec, 0xd5, 0x0e, 0x90, 0x64, + 0x69, 0x67, 0x17, 0xc1, 0x01, 0xa9, 0x63, 0x37, 0xc9, 0x28, 0xe3, 0x19, 0x33, 0xd3, 0x4e, 0xe2, + 0x3d, 0x70, 0x83, 0x03, 0xfc, 0x01, 0xee, 0xfc, 0x19, 0x8e, 0xdc, 0xb8, 0xa2, 0xf0, 0x23, 0x90, + 0xb8, 0x80, 0xba, 0xa6, 0x7b, 0x66, 0xfc, 0x15, 0x93, 0xdc, 0xa6, 0xaa, 0xab, 0xaa, 0x9f, 0x7e, + 0xea, 0xa3, 0x7b, 0xa0, 0x3a, 0x08, 0xdd, 0x0b, 0x2e, 0xc5, 0xf6, 0x20, 0x0c, 0x64, 0x40, 0x4a, + 0xae, 0x2f, 0x45, 0xe8, 0x73, 0x8f, 0x1e, 0x41, 0xd9, 0xf1, 0x7b, 0xe2, 0xea, 0x40, 0x48, 0x4e, + 0x9a, 0x50, 0xd9, 0x0f, 0xbc, 0x61, 0xdf, 0xff, 0x8c, 0x9f, 0x08, 0xaf, 0x6e, 0x35, 0xad, 0xad, + 0x32, 0xcb, 0xaa, 0x94, 0xc5, 0xb1, 0xdb, 0x17, 0x9f, 0x0f, 0xb9, 0x2f, 0x87, 0xfd, 0x7a, 0x2e, + 0xb6, 0xc8, 0xa8, 0xe8, 0xdf, 0x16, 0x94, 0x9f, 0x86, 0xbc, 0x2f, 0x30, 0xe2, 0x3a, 0x94, 0x58, + 0x70, 0x99, 0x0d, 0x97, 0xc8, 0xe4, 0x75, 0xb8, 0xef, 0xf8, 0x17, 0x22, 0x8c, 0x44, 0xdb, 0xe7, + 0x27, 0x9e, 0xe8, 0x61, 0xb8, 0x12, 0x9b, 0xd0, 0x92, 0x0d, 0x28, 0xef, 0xf3, 0xee, 0x99, 0x38, + 0x1e, 0x0d, 0x44, 0xdd, 0xc6, 0x20, 0xa9, 0x22, 0x59, 0xed, 0xb8, 0xaf, 0x44, 0x3d, 0xdf, 0xb4, + 0xb6, 0xaa, 0x2c, 0x55, 0x4c, 0xe2, 0x2d, 0x4c, 0xe1, 0x25, 0x14, 0xee, 0x31, 0xee, 0x9f, 0x26, + 0x18, 0x8a, 0x88, 0x61, 0x4c, 0x47, 0x36, 0xa1, 0xf8, 0xd4, 0x15, 0x5e, 0x2f, 0xaa, 0x2f, 0x35, + 0xed, 0xad, 0xca, 0xce, 0x83, 0x6d, 0xc3, 0xdf, 0x36, 0xea, 0x99, 0x5e, 0xa6, 0x14, 0xee, 0x3b, + 0xfd, 0x41, 0x10, 0x4a, 0x26, 0xa2, 0x41, 0xe0, 0x47, 0x82, 0xd4, 0xc0, 0x6e, 0x87, 0xa1, 0x3e, + 0xbb, 0xfa, 0xa4, 0xdf, 0x41, 0xed, 0x89, 0x17, 0x74, 0xcf, 0x5b, 0x5c, 0x72, 0x26, 0xbe, 0x1d, + 0x8a, 0x48, 0x92, 0x15, 0x28, 0x60, 0x16, 0xb4, 0x5d, 0x2c, 0x28, 0x2d, 0x32, 0xa9, 0x69, 0x8e, + 0x05, 0xa5, 0x45, 0x7f, 0xa4, 0x22, 0xcf, 0x62, 0x41, 0x69, 0x3b, 0x9e, 0xdb, 0x8d, 0x29, 0xc8, + 0xb3, 0x58, 0x20, 0x04, 0xf2, 0x2f, 0x5d, 0x71, 0xa9, 0xcf, 0x8d, 0xdf, 0xd4, 0x81, 0xe5, 0xcc, + 0xfe, 0x1a, 0xe6, 0x2a, 0x14, 0x59, 0x70, 0xe9, 0xb4, 0xa2, 0xba, 0xd5, 0xb4, 0xb7, 0xf2, 0x4c, + 0x4b, 0xc8, 0x2e, 0xa6, 0x5f, 0x2d, 0xe5, 0x70, 0x29, 0x55, 0xd0, 0x35, 0x28, 0x20, 0xd5, 0xea, + 0x94, 0xa9, 0xaf, 0xfa, 0xa4, 0xff, 0x58, 0x50, 0x3e, 0xe0, 0x57, 0x08, 0x23, 0x22, 0x1f, 0x42, + 0xa9, 0x23, 0xb9, 0xdf, 0xe3, 0x61, 0x0f, 0x8d, 0x2a, 0x3b, 0xaf, 0xa5, 0x14, 0x26, 0x66, 0xdb, + 0xc6, 0xa6, 0xed, 0xcb, 0x70, 0xc4, 0x12, 0x17, 0xb2, 0x07, 0x4b, 0xba, 0x26, 0x10, 0x43, 0x65, + 0xa7, 0x39, 0xcb, 0x3b, 0x29, 0x1b, 0xe5, 0x6c, 0x1c, 0xd6, 0x3f, 0x80, 0xea, 0x58, 0x58, 0x85, + 0xf5, 0x5c, 0x8c, 0x4c, 0x46, 0xce, 0xc5, 0x48, 0x71, 0x77, 0xc1, 0xbd, 0x61, 0xcc, 0x73, 0x9e, + 0xc5, 0xc2, 0x5e, 0xee, 0x3d, 0x6b, 0x7d, 0x0f, 0xee, 0x65, 0xa3, 0xde, 0xc6, 0x97, 0x7e, 0x0d, + 0x64, 0x3f, 0x14, 0x5c, 0x0a, 0x84, 0x77, 0x20, 0xa2, 0x88, 0x9f, 0x8a, 0xf9, 0x99, 0x8e, 0xb3, + 0x97, 0xcb, 0x66, 0x6f, 0x03, 0xca, 0x4e, 0x64, 0x0e, 0x6e, 0x63, 0x5d, 0xa6, 0x0a, 0xfa, 0x18, + 0x48, 0x4b, 0x78, 0x42, 0x0a, 0xdd, 0xbf, 0x37, 0xc4, 0xa7, 0x1d, 0x83, 0x65, 0xb1, 0x2d, 0xd9, + 0x84, 0xbc, 0x6a, 0x5d, 0x84, 0x52, 0xd9, 0xf9, 0x7f, 0xca, 0x74, 0x32, 0x27, 0x18, 0x1a, 0x50, + 0xd7, 0x04, 0xd5, 0xed, 0xbe, 0xe0, 0x80, 0x33, 0x4a, 0xd9, 0x6c, 0x65, 0x4f, 0x6e, 0x95, 0x0c, + 0x10, 0xbd, 0xd5, 0x47, 0xe6, 0xac, 0x77, 0xdd, 0x8a, 0x7e, 0xa5, 0xb5, 0xaa, 0x25, 0x0e, 0xd5, + 0x6a, 0xec, 0x83, 0xdf, 0xf3, 0x8f, 0x3c, 0x81, 0x43, 0xc5, 0x56, 0x3d, 0x14, 0xd5, 0xed, 0xa6, + 0xad, 0x62, 0xa3, 0x40, 0x77, 0xa1, 0xd8, 0xe9, 0x9e, 0x89, 0x3e, 0x27, 0x6f, 0xa8, 0x42, 0xed, + 0x89, 0x2b, 0x11, 0xe9, 0x32, 0x7f, 0x30, 0x41, 0x1f, 0x33, 0xeb, 0xf4, 0x27, 0x4b, 0xa3, 0x9f, + 0x83, 0xa8, 0x88, 0x7b, 0x47, 0xf5, 0xfc, 0xd4, 0xc4, 0x51, 0x7a, 0xa6, 0x97, 0x49, 0x1b, 0x6a, + 0x8e, 0x3f, 0x18, 0xca, 0x96, 0xf8, 0xc6, 0xf5, 0x5d, 0xe9, 0x06, 0x7e, 0x54, 0x2f, 0xa2, 0xcb, + 0x5a, 0x76, 0xeb, 0x31, 0x0b, 0x36, 0xe5, 0x42, 0x7f, 0xb0, 0xe0, 0xc1, 0x84, 0x72, 0x01, 0xae, + 0xdc, 0xcd, 0xb8, 0xde, 0x4d, 0x46, 0xa6, 0x8d, 0x86, 0x8d, 0xb9, 0x68, 0xc6, 0x27, 0xe8, 0x2f, + 0x16, 0xac, 0xcc, 0x32, 0x98, 0x89, 0xa6, 0x01, 0xf0, 0x3c, 0x74, 0xfb, 0x3c, 0x1c, 0x7d, 0x2a, + 0x46, 0xfa, 0xf6, 0xc8, 0x68, 0xc8, 0x17, 0xb0, 0x3a, 0x11, 0xeb, 0xe3, 0x6e, 0x4c, 0x51, 0x0c, + 0xea, 0xd1, 0x5c, 0x50, 0xb1, 0x1d, 0x9b, 0xe3, 0x4e, 0xff, 0xb2, 0xe0, 0xe1, 0xcc, 0xa5, 0xb4, + 0xfa, 0xac, 0x6c, 0xa1, 0x3f, 0x86, 0xda, 0x4b, 0x35, 0x18, 0x5a, 0x22, 0x92, 0xae, 0xcf, 0x95, + 0xa5, 0x2e, 0xcf, 0x29, 0x3d, 0x71, 0xa0, 0x84, 0xba, 0x03, 0x3e, 0xd0, 0x30, 0xdf, 0x5c, 0x00, + 0x73, 0xdb, 0xd8, 0xeb, 0xb9, 0x69, 0x44, 0x05, 0x06, 0xe7, 0xb8, 0xb9, 0x14, 0x50, 0x50, 0x13, + 0x71, 0xcc, 0xe1, 0x56, 0x53, 0x2d, 0x80, 0x0d, 0x33, 0x49, 0xc6, 0x90, 0xdc, 0xdc, 0x93, 0xef, + 0x03, 0xa4, 0xa6, 0xba, 0xdd, 0x6f, 0xa8, 0xcf, 0x8c, 0x31, 0x7d, 0x06, 0x1b, 0x66, 0xcc, 0xdd, + 0x62, 0x43, 0x53, 0x2d, 0xb9, 0xb4, 0x5a, 0x68, 0x1b, 0xec, 0x17, 0xcc, 0x51, 0x57, 0x1d, 0x76, + 0xab, 0x49, 0x91, 0x96, 0x94, 0xcb, 0xb3, 0x20, 0x92, 0xc6, 0x45, 0x7d, 0x2b, 0xdd, 0xf3, 0x20, + 0x94, 0x88, 0xb8, 0xca, 0xf0, 0x9b, 0x3a, 0x50, 0x3b, 0x0c, 0x7a, 0xa2, 0x23, 0xb9, 0x4c, 0x26, + 0xd1, 0x23, 0x0c, 0x8d, 0x01, 0x2b, 0x3b, 0xd5, 0xf4, 0x60, 0x2f, 0x98, 0xc3, 0x70, 0x53, 0x35, + 0xe0, 0x95, 0x83, 0x19, 0x4a, 0x28, 0xd0, 0x1f, 0x2d, 0x00, 0x13, 0x6b, 0x18, 0x2d, 0x8e, 0xf2, + 0x4e, 0xe6, 0x4e, 0x9d, 0x1e, 0x56, 0xc9, 0x12, 0xcb, 0xdc, 0xbc, 0x5b, 0x66, 0x36, 0x69, 0xd6, + 0x6b, 0xa9, 0x7d, 0xac, 0xd7, 0xe7, 0xe7, 0xf4, 0x10, 0xaa, 0xfb, 0xde, 0x30, 0x92, 0x22, 0xd4, + 0x70, 0x12, 0xcc, 0x56, 0x06, 0x33, 0xd9, 0x84, 0x25, 0x84, 0x2c, 0xa4, 0x1e, 0x01, 0x13, 0x40, + 0xcd, 0x2a, 0xed, 0x40, 0x61, 0x7e, 0xe7, 0x12, 0xc8, 0xe3, 0x73, 0x4e, 0x93, 0x8d, 0x2f, 0xb9, + 0x1a, 0xd8, 0x07, 0x6e, 0x5c, 0x1d, 0x36, 0x53, 0x9f, 0xa8, 0xe1, 0x57, 0x58, 0xbd, 0x4a, 0xc3, + 0xd5, 0x45, 0xb6, 0x1c, 0x57, 0x83, 0x9a, 0xbc, 0x77, 0xb9, 0x72, 0xcc, 0x8b, 0xc8, 0xce, 0xbc, + 0x88, 0x7e, 0xb7, 0x60, 0x99, 0x89, 0xc8, 0x7d, 0x25, 0x1c, 0x3f, 0x92, 0xe1, 0x30, 0xe9, 0xe4, + 0x4f, 0x82, 0x13, 0xa7, 0x85, 0x51, 0x6d, 0x16, 0x0b, 0x26, 0x47, 0xb9, 0xb9, 0x39, 0x7a, 0x4b, + 0xbd, 0xa1, 0x83, 0xb0, 0xa7, 0xda, 0x39, 0x08, 0x35, 0xeb, 0x13, 0x86, 0x59, 0x0b, 0xf2, 0x36, + 0x2c, 0x75, 0x82, 0x61, 0xd8, 0x4d, 0x66, 0xfd, 0x6a, 0x6a, 0x1c, 0xa3, 0x8a, 0x97, 0x99, 0x31, + 0xcb, 0xe4, 0xb4, 0xb0, 0x20, 0xa7, 0xdf, 0x5b, 0x70, 0x2f, 0x1b, 0xe3, 0x3f, 0x15, 0x6a, 0xcc, + 0x65, 0x6e, 0x26, 0x97, 0xf6, 0x2c, 0x2e, 0xf3, 0x29, 0x97, 0xe9, 0x4b, 0xa6, 0x90, 0x79, 0xc9, + 0xd0, 0x33, 0x58, 0x9b, 0x22, 0x78, 0x3f, 0xe8, 0x0f, 0x54, 0x26, 0xef, 0x4a, 0xf4, 0x0a, 0x14, + 0xda, 0x61, 0xa8, 0x29, 0x2e, 0xb3, 0x58, 0xa0, 0x5f, 0xc2, 0xc3, 0x8e, 0x90, 0x19, 0x7e, 0x33, + 0x2d, 0x7a, 0xe4, 0xf5, 0xe6, 0x9c, 0xfc, 0xc8, 0xeb, 0x29, 0x83, 0x43, 0x71, 0x39, 0x67, 0xc3, + 0x43, 0x71, 0x49, 0x77, 0xa1, 0x74, 0x1c, 0x0c, 0x02, 0x2f, 0x38, 0x1d, 0x65, 0xbb, 0xc0, 0xba, + 0xa9, 0x0b, 0x9e, 0xd4, 0x7e, 0xbd, 0x6e, 0x58, 0xbf, 0x5d, 0x37, 0xac, 0x3f, 0xae, 0x1b, 0xd6, + 0xcf, 0x7f, 0x36, 0xfe, 0x77, 0x52, 0xc4, 0x5f, 0xb0, 0xdd, 0x7f, 0x03, 0x00, 0x00, 0xff, 0xff, + 0x50, 0x2b, 0xa7, 0xda, 0x93, 0x0d, 0x00, 0x00, } diff --git a/internal/private.proto b/internal/private.proto index a97ac0e2d..bb1cadfea 100644 --- a/internal/private.proto +++ b/internal/private.proto @@ -121,6 +121,11 @@ message URI { uint32 Port = 3; } +message NodeStateMessage { + URI URI = 1; + string State = 2; +} + message NodeStatus { URI URI = 1; MaxSlices MaxSlices = 2; diff --git a/security_manager.go b/security_manager.go index 2acf82986..fc174866a 100644 --- a/security_manager.go +++ b/security_manager.go @@ -12,11 +12,7 @@ type NopSecurityManager struct { } // SetRestricted no-op. -func (sdm *NopSecurityManager) SetRestricted() { - -} +func (sdm *NopSecurityManager) SetRestricted() {} // SetNormal no-op. -func (sdm *NopSecurityManager) SetNormal() { - -} +func (sdm *NopSecurityManager) SetNormal() {} diff --git a/server.go b/server.go index c5965912a..3944fb0bf 100644 --- a/server.go +++ b/server.go @@ -151,28 +151,6 @@ func (s *Server) Open() error { s.Holder.LogOutput = s.LogOutput s.Holder.Peek() - // Start the BroadcastReceiver. - if err := s.BroadcastReceiver.Start(s); err != nil { - return fmt.Errorf("starting BroadcastReceiver: %v", err) - } - - // Open Cluster management. - if err := s.Cluster.Open(); err != nil { - return fmt.Errorf("opening Cluster: %v", err) - } - - // Open holder. - if err := s.Holder.Open(); err != nil { - return fmt.Errorf("opening Holder: %v", err) - } - - // Listen for joining nodes. - // This needs to start after the Holder has opened so that nodes can join - // the cluster without waiting for data to load on the coordinator. Before - // this starts, the joins are queued up in the Cluster.joiningURIs buffered - // channel. - s.Cluster.ListenForJoins() - // Create default HTTP client s.createDefaultClient() @@ -191,6 +169,8 @@ func (s *Server) Open() error { s.Handler.Executor = e s.Handler.LogOutput = s.LogOutput + s.Cluster.prefect = s.Handler + // Initialize Holder. s.Holder.Broadcaster = s.Broadcaster @@ -202,6 +182,29 @@ func (s *Server) Open() error { } }() + // Start the BroadcastReceiver. + if err := s.BroadcastReceiver.Start(s); err != nil { + return fmt.Errorf("starting BroadcastReceiver: %v", err) + } + + // Open Cluster management. + if err := s.Cluster.Open(); err != nil { + return fmt.Errorf("opening Cluster: %v", err) + } + + // Open holder. + if err := s.Holder.Open(); err != nil { + return fmt.Errorf("opening Holder: %v", err) + } + s.Cluster.setNodeState(NodeStateReady) + + // Listen for joining nodes. + // This needs to start after the Holder has opened so that nodes can join + // the cluster without waiting for data to load on the coordinator. Before + // this starts, the joins are queued up in the Cluster.joiningURIs buffered + // channel. + s.Cluster.ListenForJoins() + // Start background monitoring. s.wg.Add(3) go func() { defer s.wg.Done(); s.monitorAntiEntropy() }() @@ -358,6 +361,11 @@ func (s *Server) ReceiveMessage(pb proto.Message) error { } case *internal.SetCoordinatorMessage: s.Cluster.SetCoordinator(DecodeURI(obj.Old), DecodeURI(obj.New)) + case *internal.NodeStateMessage: + err := s.Cluster.ReceiveNodeState(DecodeURI(obj.URI), obj.State) + if err != nil { + return err + } } return nil diff --git a/test/handler.go b/test/handler.go index e3bee5072..6c6703f84 100644 --- a/test/handler.go +++ b/test/handler.go @@ -32,6 +32,8 @@ func NewHandler() *Handler { // Handler test messages can no-op. h.Broadcaster = pilosa.NopBroadcaster + h.SetNormal() + return h } From bd511dae8026b4ddfd02880bb403942e55e8fa1b Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Wed, 22 Nov 2017 13:42:08 -0600 Subject: [PATCH 3/9] ensure that holder opens before node is deemed ready --- cluster.go | 48 ++++++++++++++++++++++++++++-------------------- holder.go | 7 +++++++ server.go | 8 +++++--- test/cluster.go | 9 +++++++-- 4 files changed, 47 insertions(+), 25 deletions(-) diff --git a/cluster.go b/cluster.go index 171615b4d..7ea161215 100644 --- a/cluster.go +++ b/cluster.go @@ -286,13 +286,13 @@ func (c *Cluster) setState(state string) { // - ClusterStateStarting } + c.logger().Printf("Change cluster state from %s to %s", c.State, state) c.State = state } -func (c *Cluster) setNodeState(state string) { +func (c *Cluster) SetNodeState(state string) error { if c.IsCoordinator() { - c.Topology.nodeStates[c.URI] = state - return + return c.ReceiveNodeState(c.URI, state) } // Send node state to coordinator. @@ -305,8 +305,10 @@ func (c *Cluster) setNodeState(state string) { URI: c.Coordinator, } if err := c.Broadcaster.SendTo(node, ns); err != nil { - c.logger().Printf("sending node state error: err=%s", err) + return fmt.Errorf("sending node state error: err=%s", err) } + + return nil } func (c *Cluster) ReceiveNodeState(uri URI, state string) error { @@ -691,18 +693,15 @@ func (c *Cluster) Open() error { // Only the coordinator needs to consider the .topology file. if c.IsCoordinator() { - state, err := c.considerTopology() + err := c.considerTopology() if err != nil { return fmt.Errorf("considerTopology: %v", err) } - // Add the local node to the cluster and update state. - c.AddNode(c.URI) - c.setState(state) - } else { - // Add the local node to the cluster. - c.AddNode(c.URI) } + // Add the local node to the cluster. + c.AddNode(c.URI) + // Start the EventReceiver. if err := c.EventReceiver.Start(c); err != nil { return fmt.Errorf("starting EventReceiver: %v", err) @@ -921,6 +920,10 @@ func (c *Cluster) CompleteCurrentJob(state string) error { // FollowResizeInstruction is run by any node that receives a ResizeInstruction. func (c *Cluster) FollowResizeInstruction(instr *internal.ResizeInstruction) error { go func() { + + // Make sure the holder has opened. + <-c.Holder.opened + // Prepare the return message. complete := &internal.ResizeInstructionComplete{ JobID: instr.JobID, @@ -1273,6 +1276,11 @@ func (c *Cluster) loadTopology() error { // saveTopology writes the current topology to disk. func (c *Cluster) saveTopology() error { + + if err := os.MkdirAll(c.Path, 0777); err != nil { + return err + } + if buf, err := proto.Marshal(encodeTopology(c.Topology)); err != nil { return err } else if err := ioutil.WriteFile(filepath.Join(c.Path, ".topology"), buf, 0666); err != nil { @@ -1301,25 +1309,25 @@ func decodeTopology(topology *internal.Topology) (*Topology, error) { return t, nil } -func (c *Cluster) considerTopology() (string, error) { - // If there is no .topology file, it's safe to go to state NORMAL. +func (c *Cluster) considerTopology() error { + // If there is no .topology file, it's safe to proceed. if len(c.Topology.NodeSet) == 0 { - return ClusterStateNormal, nil + return nil } // The local node (coordinator) must be in the .topology. if !c.Topology.ContainsURI(c.Coordinator) { - return "", fmt.Errorf("coordinator %s is not in topology: %v", c.Coordinator, c.Topology.NodeSet) + return fmt.Errorf("coordinator %s is not in topology: %v", c.Coordinator, c.Topology.NodeSet) } - // If local node is the only thing in .topology, continue to state NORMAL. - if len(c.Topology.NodeSet) == 1 { - return ClusterStateNormal, nil - } + // If local node is the only thing in .topology, continue. + //if len(c.Topology.NodeSet) == 1 { + // return nil + //} // Keep the cluster in state "STARTING" until hearing from all nodes. // Topology contains 2+ hosts. - return ClusterStateStarting, nil + return nil } // ReceiveEvent represents an implementation of EventHandler. diff --git a/holder.go b/holder.go index 0e3d3953c..cd700e586 100644 --- a/holder.go +++ b/holder.go @@ -46,6 +46,9 @@ type Holder struct { indexes map[string]*Index hasData bool + // opened channel is closed once Open() completes. + opened chan struct{} + Broadcaster Broadcaster // Close management wg sync.WaitGroup @@ -69,6 +72,8 @@ func NewHolder() *Holder { indexes: make(map[string]*Index), closing: make(chan struct{}, 0), + opened: make(chan struct{}), + Broadcaster: NopBroadcaster, Stats: NopStatsClient, @@ -156,6 +161,8 @@ func (h *Holder) Open() error { go func() { defer h.wg.Done(); h.monitorCacheFlush() }() h.Stats.Open() + + close(h.opened) return nil } diff --git a/server.go b/server.go index 3944fb0bf..b677147bd 100644 --- a/server.go +++ b/server.go @@ -196,13 +196,15 @@ func (s *Server) Open() error { if err := s.Holder.Open(); err != nil { return fmt.Errorf("opening Holder: %v", err) } - s.Cluster.setNodeState(NodeStateReady) + if err := s.Cluster.SetNodeState(NodeStateReady); err != nil { + return fmt.Errorf("setting nodeState: %v", err) + } // Listen for joining nodes. // This needs to start after the Holder has opened so that nodes can join // the cluster without waiting for data to load on the coordinator. Before - // this starts, the joins are queued up in the Cluster.joiningURIs buffered - // channel. + // this starts, the joins are queued up in the Cluster.joiningLeavingNodes + // buffered channel. s.Cluster.ListenForJoins() // Start background monitoring. diff --git a/test/cluster.go b/test/cluster.go index 82c45aef6..449f50580 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -250,8 +250,13 @@ func (t *TestCluster) SetState(state string) { // Open opens all clusters in the test cluster. func (t *TestCluster) Open() error { for _, c := range t.Clusters { - err := c.Open() - if err != nil { + if err := c.Open(); err != nil { + return err + } + if err := c.Holder.Open(); err != nil { + return err + } + if err := c.SetNodeState(pilosa.NodeStateReady); err != nil { return err } } From f2c32f8ec9a5ca8cd375dfdeb40b3be7eb37b5bf Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Wed, 22 Nov 2017 15:04:26 -0600 Subject: [PATCH 4/9] add logging support to ResizeJob --- cluster.go | 21 +++++++++++++++++---- 1 file changed, 17 insertions(+), 4 deletions(-) diff --git a/cluster.go b/cluster.go index 7ea161215..39bf5afff 100644 --- a/cluster.go +++ b/cluster.go @@ -831,6 +831,7 @@ func (c *Cluster) listenForJoins() { // added/removed. It also saves a reference to the ResizeJob in the `jobs` map // for future lookup by JobID. func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*ResizeJob, error) { + c.logger().Printf("generateResizeJob: %v", nodeAction) c.mu.Lock() defer c.mu.Unlock() @@ -838,6 +839,7 @@ func (c *Cluster) generateResizeJob(nodeAction nodeAction) (*ResizeJob, error) { if err != nil { return nil, err } + c.logger().Printf("generated ResizeJob: %d", j.ID) // Save job in jobs map for future reference. c.jobs[j.ID] = j @@ -1049,6 +1051,14 @@ type ResizeJob struct { mu sync.RWMutex state string + + // The writer for any logging. + LogOutput io.Writer +} + +// logger returns a logger for the resize job. +func (j *ResizeJob) logger() *log.Logger { + return log.New(j.LogOutput, "", log.LstdFlags) } // NewResizeJob returns a new instance of ResizeJob. @@ -1076,10 +1086,11 @@ func NewResizeJob(existingURIs []URI, uri URI, action string) *ResizeJob { } return &ResizeJob{ - ID: rand.Int63(), - URIs: uris, - action: action, - result: make(chan string), + ID: rand.Int63(), + URIs: uris, + action: action, + result: make(chan string), + LogOutput: os.Stderr, } } @@ -1141,6 +1152,7 @@ func (j *ResizeJob) urisArePending() bool { } func (j *ResizeJob) distributeResizeInstructions() error { + j.logger().Printf("distributeResizeInstructions for job %d", j.ID) // Loop through the ResizeInstructions in ResizeJob and send to each host. for _, instr := range j.Instructions { // Because the node may not be in the cluster yet, create @@ -1148,6 +1160,7 @@ func (j *ResizeJob) distributeResizeInstructions() error { node := &Node{ URI: decodeURI(instr.URI), } + j.logger().Printf("send resize instructions: %v", instr) if err := j.Broadcaster.SendTo(node, instr); err != nil { return err } From 013cd0cd95d3690fe3f956f0aa097551455be61a Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Wed, 22 Nov 2017 17:08:26 -0600 Subject: [PATCH 5/9] refactor resize instruction logic to support multi-index. wait to open holder on non-coordinator nodes. --- cluster.go | 42 ++++++++++++++++++++++++++---------------- server.go | 20 ++++++++++++++++++++ 2 files changed, 46 insertions(+), 16 deletions(-) diff --git a/cluster.go b/cluster.go index 39bf5afff..d66883f5e 100644 --- a/cluster.go +++ b/cluster.go @@ -876,33 +876,43 @@ func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*ResizeJob, pbSchema := c.Holder.EncodeSchema() - // Add to the ResizeJob the instructions for each index. + // multiIndex is a map of sources for each node in toCluster. + // Initialize the map with all the nodes in toCluster. + multiIndex := make(map[URI][]*internal.ResizeSource) + for _, n := range toCluster.Nodes { + multiIndex[n.URI] = nil + } + + // Add to m the instructions for each index. for _, idx := range c.Holder.Indexes() { - // fragSources is map[URI][]*internal.ResizeSource. fragSources, err := c.fragSources(toCluster, idx) if err != nil { return nil, err } for u, sources := range fragSources { - // If a host doesn't need to request data, mark it as complete. - if len(sources) == 0 { - j.URIs[u] = true - continue + for _, src := range sources { + multiIndex[u] = append(multiIndex[u], src) } - // 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: u.Encode(), - Coordinator: encodeURI(c.Coordinator), - Sources: sources, - Schema: pbSchema, // Include the schema to ensure it's in sync on the receiving node. - } - j.Instructions = append(j.Instructions, instr) } } + for u, sources := range multiIndex { + // If a host doesn't need to request data, mark it as complete. + if len(sources) == 0 { + j.URIs[u] = true + continue + } + instr := &internal.ResizeInstruction{ + JobID: j.ID, + URI: u.Encode(), + Coordinator: encodeURI(c.Coordinator), + Sources: sources, + Schema: pbSchema, // Include the schema to ensure it's in sync on the receiving node. + } + j.Instructions = append(j.Instructions, instr) + } + return j, nil } diff --git a/server.go b/server.go index b677147bd..7328eac80 100644 --- a/server.go +++ b/server.go @@ -51,6 +51,11 @@ type Server struct { wg sync.WaitGroup closing chan struct{} + // joining is held open until this node + // receives ClusterStatus from the coordinator. + joining chan struct{} + joined bool + // Data storage and HTTP interface. Holder *Holder Handler *Handler @@ -84,6 +89,7 @@ type Server struct { func NewServer() *Server { s := &Server{ closing: make(chan struct{}), + joining: make(chan struct{}), Holder: NewHolder(), Handler: NewHandler(), @@ -192,6 +198,12 @@ func (s *Server) Open() error { return fmt.Errorf("opening Cluster: %v", err) } + // If not coordinator then wait for ClusterStatus from coordinator. + if !s.Cluster.IsCoordinator() { + s.Logger().Printf("wait for joining to complete") + <-s.joining + } + // Open holder. if err := s.Holder.Open(); err != nil { return fmt.Errorf("opening Holder: %v", err) @@ -216,6 +228,13 @@ func (s *Server) Open() error { return nil } +func (s *Server) markAsJoined() { + if !s.joined { + s.joined = true + close(s.joining) + } +} + // Close closes the server and waits for it to shutdown. func (s *Server) Close() error { // Notify goroutines to stop. @@ -351,6 +370,7 @@ func (s *Server) ReceiveMessage(pb proto.Message) error { if err != nil { return err } + s.markAsJoined() case *internal.ResizeInstruction: err := s.Cluster.FollowResizeInstruction(obj) if err != nil { From ae63adfaac349b0de7da156966c9e2891e4db380 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 27 Nov 2017 10:32:49 -0600 Subject: [PATCH 6/9] let remote nodes know that its safe to open Holder when launching based on existing taxonomy --- cluster.go | 28 +++++++++++++++++----------- 1 file changed, 17 insertions(+), 11 deletions(-) diff --git a/cluster.go b/cluster.go index d66883f5e..9f30c0f94 100644 --- a/cluster.go +++ b/cluster.go @@ -301,10 +301,7 @@ func (c *Cluster) SetNodeState(state string) error { State: state, } - node := &Node{ - URI: c.Coordinator, - } - if err := c.Broadcaster.SendTo(node, ns); err != nil { + if err := c.sendTo(c.Coordinator, ns); err != nil { return fmt.Errorf("sending node state error: err=%s", err) } @@ -779,6 +776,14 @@ func (c *Cluster) setStateAndBroadcast(state string) error { return c.Broadcaster.SendSync(c.Status()) } +func (c *Cluster) sendTo(to URI, msg proto.Message) error { + node := &Node{URI: to} + if err := c.Broadcaster.SendTo(node, msg); err != nil { + return err + } + return nil +} + // ListenForJoins handles cluster-resize events. func (c *Cluster) ListenForJoins() { c.wg.Add(1) @@ -876,14 +881,14 @@ func (c *Cluster) generateResizeJobByAction(nodeAction nodeAction) (*ResizeJob, pbSchema := c.Holder.EncodeSchema() - // multiIndex is a map of sources for each node in toCluster. - // Initialize the map with all the nodes in toCluster. + // multiIndex is a map of sources initialized with all the nodes in toCluster. multiIndex := make(map[URI][]*internal.ResizeSource) + for _, n := range toCluster.Nodes { multiIndex[n.URI] = nil } - // Add to m the instructions for each index. + // Add to multiIndex the instructions for each index. for _, idx := range c.Holder.Indexes() { fragSources, err := c.fragSources(toCluster, idx) if err != nil { @@ -1002,10 +1007,7 @@ func (c *Cluster) FollowResizeInstruction(instr *internal.ResizeInstruction) err complete.Error = err.Error() } - node := &Node{ - URI: decodeURI(instr.Coordinator), - } - if err := c.Broadcaster.SendTo(node, complete); err != nil { + if err := c.sendTo(decodeURI(instr.Coordinator), complete); err != nil { c.logger().Printf("sending resizeInstructionComplete error: err=%s", err) } }() @@ -1401,6 +1403,10 @@ func (c *Cluster) nodeJoin(uri URI) error { if c.haveTopologyAgreement() && c.allNodesReady() { return c.setStateAndBroadcast(ClusterStateNormal) + } else { + // Send the status to the remote node. This lets the remote node + // know that it can proceed with opening its Holder. + return c.sendTo(uri, c.Status()) } return nil From ccdd6262c59a8c3ad107fea2fc337956044ef4c3 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Tue, 28 Nov 2017 12:43:49 -0600 Subject: [PATCH 7/9] add error logging for non-topology nodeJoin --- cluster.go | 4 +++- gossip/gossip.go | 4 ++-- server.go | 21 ++++----------------- 3 files changed, 9 insertions(+), 20 deletions(-) diff --git a/cluster.go b/cluster.go index 9f30c0f94..6d02d352b 100644 --- a/cluster.go +++ b/cluster.go @@ -1382,7 +1382,9 @@ 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) + err := fmt.Sprintf("host is not in topology: %v", uri) + c.logger().Print(err) + return errors.New(err) } if err := c.AddNode(uri); err != nil { diff --git a/gossip/gossip.go b/gossip/gossip.go index 2c5c85faf..7e89e17c1 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. g.config.memberlistConfig.Delegate = g g.config.memberlistConfig.SecretKey = secretKey g.config.memberlistConfig.Events = server.Cluster.EventReceiver.(memberlist.EventDelegate) @@ -336,7 +336,7 @@ func (g *GossipEventReceiver) listen() { Event: nodeEventType, URI: *uri, } - g.eventHandler.ReceiveEvent(ne) + _ = g.eventHandler.ReceiveEvent(ne) // TODO: don't swallow this error } } diff --git a/server.go b/server.go index 7328eac80..8a4ad38d5 100644 --- a/server.go +++ b/server.go @@ -432,6 +432,10 @@ func (s *Server) ClusterStatus() (proto.Message, error) { // HandleRemoteStatus receives incoming NodeStatus from remote nodes. func (s *Server) HandleRemoteStatus(pb proto.Message) error { + // Ignore NodeStatus messages until the cluster is in a Normal state. + if s.Cluster.State != ClusterStateNormal { + return nil + } return s.mergeRemoteStatus(pb.(*internal.NodeStatus)) } @@ -441,23 +445,6 @@ 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 From 8228ef238203eaebc0a41ad62ad38d390bc44deb Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Tue, 28 Nov 2017 16:40:05 -0600 Subject: [PATCH 8/9] Add a Cluster.Static override for tests to treat static nodes as Coordinator. Refactor the Server.joining channel to be in Cluster instead. --- cluster.go | 32 +++++++++++++++++++++++++++++++- server.go | 21 +-------------------- server/server.go | 14 +++++++++++++- 3 files changed, 45 insertions(+), 22 deletions(-) diff --git a/cluster.go b/cluster.go index 6d02d352b..50996d6f1 100644 --- a/cluster.go +++ b/cluster.go @@ -166,6 +166,7 @@ type Cluster struct { Topology *Topology // Required for cluster Resize. + Static bool // Static is primarily used for testing in a non-gossip environment. State string Coordinator URI Holder *Holder @@ -173,6 +174,11 @@ type Cluster struct { joiningLeavingNodes chan nodeAction + // joining is held open until this node + // receives ClusterStatus from the coordinator. + joining chan struct{} + joined bool + mu sync.RWMutex jobs map[int64]*ResizeJob currentJob *ResizeJob @@ -197,6 +203,7 @@ func NewCluster() *Cluster { joiningLeavingNodes: make(chan nodeAction, 10), // buffered channel jobs: make(map[int64]*ResizeJob), closing: make(chan struct{}), + joining: make(chan struct{}), LogOutput: os.Stderr, prefect: &NopSecurityManager{}, @@ -210,7 +217,7 @@ func (c *Cluster) logger() *log.Logger { // IsCoordinator is true if this node is the coordinator. func (c *Cluster) IsCoordinator() bool { - return c.Coordinator == c.URI + return c.Static || c.Coordinator == c.URI } // SetCoordinator updates the Coordinator to new if it is @@ -709,6 +716,12 @@ func (c *Cluster) Open() error { return fmt.Errorf("opening MemberSet: %v", err) } + // If not coordinator then wait for ClusterStatus from coordinator. + if !c.IsCoordinator() { + c.logger().Printf("wait for joining to complete") + <-c.joining + } + return nil } @@ -720,15 +733,28 @@ func (c *Cluster) Close() error { return nil } +func (c *Cluster) MarkAsJoined() { + if !c.joined { + c.joined = true + close(c.joining) + } +} + func (c *Cluster) needTopologyAgreement() bool { return c.State == ClusterStateStarting && !URISlicesAreEqual(c.Topology.NodeSet, c.NodeSet()) } func (c *Cluster) haveTopologyAgreement() bool { + if c.Static { + return true + } return URISlicesAreEqual(c.Topology.NodeSet, c.NodeSet()) } func (c *Cluster) allNodesReady() bool { + if c.Static { + return true + } for _, uri := range c.Topology.NodeSet { if c.Topology.nodeStates[uri] != NodeStateReady { return false @@ -1335,6 +1361,10 @@ func decodeTopology(topology *internal.Topology) (*Topology, error) { } func (c *Cluster) considerTopology() error { + if c.Static { + return nil + } + // If there is no .topology file, it's safe to proceed. if len(c.Topology.NodeSet) == 0 { return nil diff --git a/server.go b/server.go index 8a4ad38d5..37270b0f1 100644 --- a/server.go +++ b/server.go @@ -51,11 +51,6 @@ type Server struct { wg sync.WaitGroup closing chan struct{} - // joining is held open until this node - // receives ClusterStatus from the coordinator. - joining chan struct{} - joined bool - // Data storage and HTTP interface. Holder *Holder Handler *Handler @@ -89,7 +84,6 @@ type Server struct { func NewServer() *Server { s := &Server{ closing: make(chan struct{}), - joining: make(chan struct{}), Holder: NewHolder(), Handler: NewHandler(), @@ -198,12 +192,6 @@ func (s *Server) Open() error { return fmt.Errorf("opening Cluster: %v", err) } - // If not coordinator then wait for ClusterStatus from coordinator. - if !s.Cluster.IsCoordinator() { - s.Logger().Printf("wait for joining to complete") - <-s.joining - } - // Open holder. if err := s.Holder.Open(); err != nil { return fmt.Errorf("opening Holder: %v", err) @@ -228,13 +216,6 @@ func (s *Server) Open() error { return nil } -func (s *Server) markAsJoined() { - if !s.joined { - s.joined = true - close(s.joining) - } -} - // Close closes the server and waits for it to shutdown. func (s *Server) Close() error { // Notify goroutines to stop. @@ -370,7 +351,7 @@ func (s *Server) ReceiveMessage(pb proto.Message) error { if err != nil { return err } - s.markAsJoined() + s.Cluster.MarkAsJoined() case *internal.ResizeInstruction: err := s.Cluster.FollowResizeInstruction(obj) if err != nil { diff --git a/server/server.go b/server/server.go index b4d067683..5237df085 100644 --- a/server/server.go +++ b/server/server.go @@ -63,7 +63,7 @@ type Command struct { // Standard input/output *pilosa.CmdIO - // running will be closed once Command.Run is finished. + // Started will be closed once Command.Run is finished. Started chan struct{} // Done will be closed when Command.Close() is called Done chan struct{} @@ -220,6 +220,18 @@ func (m *Command) SetupServer() error { m.Server.Broadcaster = gossipMemberSet m.Server.BroadcastReceiver = gossipMemberSet case pilosa.ClusterStatic, pilosa.ClusterNone: + + m.Server.Cluster.Static = true + for _, address := range m.Config.Cluster.Hosts { + uri, err := pilosa.NewURIFromAddress(address) + if err != nil { + return err + } + cluster.Nodes = append(cluster.Nodes, &pilosa.Node{ + URI: *uri, + }) + } + m.Server.Broadcaster = pilosa.NopBroadcaster m.Server.Cluster.MemberSet = pilosa.NewStaticMemberSet() m.Server.BroadcastReceiver = pilosa.NopBroadcastReceiver From 91c7abfa9a01a4a2ca71b35fe55eea25318e6703 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Tue, 28 Nov 2017 16:47:40 -0600 Subject: [PATCH 9/9] remove outdated TODO --- handler.go | 2 -- 1 file changed, 2 deletions(-) diff --git a/handler.go b/handler.go index 1ef7be7f9..11b1f0cb9 100644 --- a/handler.go +++ b/handler.go @@ -2010,8 +2010,6 @@ func (h *Handler) handlePostClusterResizeRemoveNode(w http.ResponseWriter, r *ht // 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