From f17a399a37fbbc1f4e64788328282db528cd8880 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 7 Dec 2017 08:33:26 -0600 Subject: [PATCH 1/4] adjust to single http client --- cluster.go | 6 +++++- 1 file changed, 5 insertions(+), 1 deletion(-) diff --git a/cluster.go b/cluster.go index c44f35d55..f6966a333 100644 --- a/cluster.go +++ b/cluster.go @@ -24,6 +24,7 @@ import ( "io/ioutil" "log" "math/rand" + "net/http" "os" "path/filepath" "sort" @@ -190,6 +191,9 @@ type Cluster struct { // The writer for any logging. LogOutput io.Writer + + // + RemoteClient *http.Client } // NewCluster returns a new instance of Cluster with defaults. @@ -1031,7 +1035,7 @@ func (c *Cluster) FollowResizeInstruction(instr *internal.ResizeInstruction) err } // Create a client for calling remote nodes. - client := NewInternalHTTPClientFromURI(&c.URI, nil) // TODO: ClientOptions + client := NewInternalHTTPClientFromURI(&c.URI, GetHTTPClient(nil)) // TODO: ClientOptions // Request each source file in ResizeSources. for _, src := range instr.Sources { From c14bc2904114f7d0fd622842b442c9e525a1b5f3 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 7 Dec 2017 11:17:37 -0600 Subject: [PATCH 2/4] added logging for node membership --- cluster.go | 11 ++++++++++- holder.go | 2 ++ server/server.go | 1 + 3 files changed, 13 insertions(+), 1 deletion(-) diff --git a/cluster.go b/cluster.go index f6966a333..368cddba4 100644 --- a/cluster.go +++ b/cluster.go @@ -336,6 +336,7 @@ func (c *Cluster) SetNodeState(state string) error { State: state, } + c.logger().Printf("Sending State %s (%s)", state, c.Coordinator.String()) if err := c.sendTo(c.Coordinator, ns); err != nil { return fmt.Errorf("sending node state error: err=%s", err) } @@ -347,6 +348,8 @@ func (c *Cluster) SetNodeState(state string) error { // Coordinator to keep track of, during startup, which nodes have // finished opening their Holder. func (c *Cluster) ReceiveNodeState(uri URI, state string) error { + + c.logger().Printf("Receiving State %s (%s)", state, uri.String()) if !c.IsCoordinator() { return nil } @@ -357,9 +360,11 @@ func (c *Cluster) ReceiveNodeState(uri URI, state string) error { } c.Topology.nodeStates[uri] = state + c.logger().Printf("Receiving State %s (%s)", state, uri.String()) // Set cluster state to NORMAL. if c.haveTopologyAgreement() && c.allNodesReady() { + c.logger().Printf("Broadcasting ClusterStateNormal") return c.setStateAndBroadcast(ClusterStateNormal) } @@ -800,6 +805,8 @@ func (c *Cluster) haveTopologyAgreement() bool { if c.Static { return true } + c.logger().Printf("haveTopologyAgreement") + c.logger().Printf(" (%v)(%v)", c.Topology.NodeSet, c.NodeSet()) return URISlicesAreEqual(c.Topology.NodeSet, c.NodeSet()) } @@ -807,7 +814,9 @@ func (c *Cluster) allNodesReady() bool { if c.Static { return true } + c.logger().Printf("allNodesReady") for _, uri := range c.Topology.NodeSet { + c.logger().Printf("allNodesReady: %s,%s", uri.String(), c.Topology.nodeStates[uri]) if c.Topology.nodeStates[uri] != NodeStateReady { return false } @@ -1035,7 +1044,7 @@ func (c *Cluster) FollowResizeInstruction(instr *internal.ResizeInstruction) err } // Create a client for calling remote nodes. - client := NewInternalHTTPClientFromURI(&c.URI, GetHTTPClient(nil)) // TODO: ClientOptions + client := NewInternalHTTPClientFromURI(&c.URI, c.RemoteClient) // TODO: ClientOptions // Request each source file in ResizeSources. for _, src := range instr.Sources { diff --git a/holder.go b/holder.go index 122dafd7a..99d9d079d 100644 --- a/holder.go +++ b/holder.go @@ -133,6 +133,7 @@ func (h *Holder) Open() error { return err } + h.logger().Printf("Holder Start") for _, fi := range fis { if !fi.IsDir() { continue @@ -156,6 +157,7 @@ func (h *Holder) Open() error { } h.indexes[index.Name()] = index } + h.logger().Printf("Holder Complete") // Periodically flush cache. h.wg.Add(1) diff --git a/server/server.go b/server/server.go index 9af9a1947..e7396dfb1 100644 --- a/server/server.go +++ b/server/server.go @@ -181,6 +181,7 @@ func (m *Command) SetupServer() error { c := pilosa.GetHTTPClient(TLSConfig) m.Server.RemoteClient = c m.Server.Handler.RemoteClient = c + m.Server.Cluster.RemoteClient = c // Set the coordinator node. curi, err := pilosa.AddressWithDefaults(m.Config.Cluster.Coordinator) From 9565db43a83f25d53c853e01ede26effe642e6aa Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 7 Dec 2017 14:11:49 -0600 Subject: [PATCH 3/4] add sort order to URI addition --- cluster.go | 14 ++++++++++---- 1 file changed, 10 insertions(+), 4 deletions(-) diff --git a/cluster.go b/cluster.go index 368cddba4..a17323783 100644 --- a/cluster.go +++ b/cluster.go @@ -805,8 +805,6 @@ func (c *Cluster) haveTopologyAgreement() bool { if c.Static { return true } - c.logger().Printf("haveTopologyAgreement") - c.logger().Printf(" (%v)(%v)", c.Topology.NodeSet, c.NodeSet()) return URISlicesAreEqual(c.Topology.NodeSet, c.NodeSet()) } @@ -814,9 +812,7 @@ func (c *Cluster) allNodesReady() bool { if c.Static { return true } - c.logger().Printf("allNodesReady") for _, uri := range c.Topology.NodeSet { - c.logger().Printf("allNodesReady: %s,%s", uri.String(), c.Topology.nodeStates[uri]) if c.Topology.nodeStates[uri] != NodeStateReady { return false } @@ -1343,6 +1339,12 @@ func (t *Topology) AddURI(uri URI) bool { return false } t.NodeSet = append(t.NodeSet, uri) + + sort.Slice(t.NodeSet, + func(i, j int) bool { + return t.NodeSet[i].String() < t.NodeSet[j].String() + }) + return true } @@ -1422,6 +1424,10 @@ func decodeTopology(topology *internal.Topology) (*Topology, error) { t := NewTopology() t.NodeSet = decodeURIs(topology.NodeSet) + sort.Slice(t.NodeSet, + func(i, j int) bool { + return t.NodeSet[i].String() < t.NodeSet[j].String() + }) return t, nil } From e6adb361b90681bc766f883bd38d851c28067119 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Thu, 7 Dec 2017 15:09:42 -0600 Subject: [PATCH 4/4] cleanup logging --- cluster.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/cluster.go b/cluster.go index a17323783..1fddefcf3 100644 --- a/cluster.go +++ b/cluster.go @@ -336,7 +336,7 @@ func (c *Cluster) SetNodeState(state string) error { State: state, } - c.logger().Printf("Sending State %s (%s)", state, c.Coordinator.String()) + c.logger().Printf("Sending State %s (%s)", state, c.Coordinator) if err := c.sendTo(c.Coordinator, ns); err != nil { return fmt.Errorf("sending node state error: err=%s", err) } @@ -360,7 +360,7 @@ func (c *Cluster) ReceiveNodeState(uri URI, state string) error { } c.Topology.nodeStates[uri] = state - c.logger().Printf("Receiving State %s (%s)", state, uri.String()) + c.logger().Printf("Receiving State %s (%s)", state, uri) // Set cluster state to NORMAL. if c.haveTopologyAgreement() && c.allNodesReady() {