From 9d870197625a2ded7c415c512cb0deea9cb5c84f Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Fri, 10 Nov 2017 09:10:15 -0600 Subject: [PATCH 1/3] WIP: don't block joining nodes while coordinator loads data. Implements a Holder.Peek() function, and breaks out Cluster.ListenForJoins() into a separate method that can be started after the Holder finishes loading data. TODO: - [ ] test Holder.Peek() --- cluster.go | 14 ++++++++------ holder.go | 33 ++++++++++++++++++++++++++++++++- server.go | 20 +++++++++++++++++--- 3 files changed, 57 insertions(+), 10 deletions(-) diff --git a/cluster.go b/cluster.go index 32f1dee12..204e39c26 100644 --- a/cluster.go +++ b/cluster.go @@ -546,10 +546,6 @@ func (c *Cluster) Open() error { return fmt.Errorf("opening MemberSet: %v", err) } - // Listen for cluster-resize events. - c.wg.Add(1) - go func() { defer c.wg.Done(); c.listenForJoins() }() - return nil } @@ -604,6 +600,12 @@ func (c *Cluster) setStateAndBroadcast(state string) error { return c.Broadcaster.SendSync(c.Status()) } +// ListenForJoins handles cluster-resize events. +func (c *Cluster) ListenForJoins() { + c.wg.Add(1) + go func() { defer c.wg.Done(); c.listenForJoins() }() +} + func (c *Cluster) listenForJoins() { var uriJoined bool @@ -634,8 +636,8 @@ func (c *Cluster) listenForJoins() { select { case <-c.closing: return - case host := <-c.joiningURIs: - err := c.handleJoiningHost(host) + case uri := <-c.joiningURIs: + err := c.handleJoiningHost(uri) if err != nil { c.logger().Printf("handleJoiningHost error: err=%s", err) continue diff --git a/holder.go b/holder.go index 07c030135..bbda143d6 100644 --- a/holder.go +++ b/holder.go @@ -44,6 +44,7 @@ type Holder struct { // Indexes by name. indexes map[string]*Index + hasData bool Broadcaster Broadcaster // Close management @@ -77,6 +78,36 @@ func NewHolder() *Holder { } } +// Peek reads the root data directory for the holder +// without actually loading any data into memory. +func (h *Holder) Peek() error { + if err := os.MkdirAll(h.Path, 0777); err != nil { + return err + } + + // Open path to read all index directories. + f, err := os.Open(h.Path) + if err != nil { + return err + } + defer f.Close() + + fis, err := f.Readdir(0) + if err != nil { + return err + } + + for _, fi := range fis { + if !fi.IsDir() { + continue + } + h.hasData = true + break + } + + return nil +} + // Open initializes the root data directory for the holder. func (h *Holder) Open() error { h.setFileLimit() @@ -149,7 +180,7 @@ func (h *Holder) Close() error { // This is used to determine if the rebalancing of data is necessary // when a node joins the cluster. func (h *Holder) HasData() bool { - return len(h.indexes) > 0 + return h.hasData || len(h.indexes) > 0 } // MaxSlices returns MaxSlice map for all indexes. diff --git a/server.go b/server.go index ac5c5d06a..06c0a81e7 100644 --- a/server.go +++ b/server.go @@ -145,10 +145,12 @@ func (s *Server) Open() error { } } - // Open holder. + // Peek at the holder to determine if there is data on disk. + // Don't actually load the data until after the Cluster + // management starts. s.Holder.LogOutput = s.LogOutput - if err := s.Holder.Open(); err != nil { - return fmt.Errorf("opening Holder: %v", err) + if err := s.Holder.Peek(); err != nil { + return fmt.Errorf("peeking at the Holder: %v", err) } // Start the BroadcastReceiver. @@ -161,6 +163,18 @@ func (s *Server) Open() error { 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() From 453bbc996c0ad3342a9683ae92b1eceed52618f3 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Fri, 10 Nov 2017 10:38:49 -0600 Subject: [PATCH 2/3] adjust Holder.Peek() and add tests for it --- holder.go | 13 ++++++------ holder_test.go | 54 ++++++++++++++++++++++++++++++++++++++++++++++++++ server.go | 4 +--- 3 files changed, 61 insertions(+), 10 deletions(-) diff --git a/holder.go b/holder.go index bbda143d6..0e3d3953c 100644 --- a/holder.go +++ b/holder.go @@ -80,21 +80,20 @@ func NewHolder() *Holder { // Peek reads the root data directory for the holder // without actually loading any data into memory. -func (h *Holder) Peek() error { - if err := os.MkdirAll(h.Path, 0777); err != nil { - return err - } +// HasData is returned, and h.hasData is set. +func (h *Holder) Peek() bool { + h.hasData = false // Open path to read all index directories. f, err := os.Open(h.Path) if err != nil { - return err + return false } defer f.Close() fis, err := f.Readdir(0) if err != nil { - return err + return false } for _, fi := range fis { @@ -105,7 +104,7 @@ func (h *Holder) Peek() error { break } - return nil + return h.hasData } // Open initializes the root data directory for the holder. diff --git a/holder_test.go b/holder_test.go index 0427c577e..f182d3496 100644 --- a/holder_test.go +++ b/holder_test.go @@ -266,6 +266,60 @@ func TestHolder_Open(t *testing.T) { }) } +func TestHolder_HasData(t *testing.T) { + t.Run("IndexDirectory", func(t *testing.T) { + h := test.MustOpenHolder() + defer h.Close() + + if h.HasData() { + t.Fatal("expected HasData to return false") + } + + if _, err := h.CreateIndex("test", pilosa.IndexOptions{}); err != nil { + t.Fatal(err) + } + + if !h.HasData() { + t.Fatal("expected HasData to return true") + } + }) + + t.Run("Peek", func(t *testing.T) { + h := test.NewHolder() + + if hasData := h.Peek(); hasData != false { + t.Fatal("expected Peek to return false") + } else if h.HasData() { + t.Fatal("expected HasData to return false") + } + + // Create an index directory to indicate data exists. + if err := os.Mkdir(h.IndexPath("test"), 0777); err != nil { + t.Fatal(err) + } + + if hasData := h.Peek(); hasData != true { + t.Fatal("expected Peek to return true") + } else if !h.HasData() { + t.Fatal("expected HasData to return true") + } + }) + + t.Run("Peek at missing directory", func(t *testing.T) { + h := test.NewHolder() + + // Ensure that hasData is false when trying to peek into + // a directory that doesn't exist. + h.Path = "bad-path" + + if hasData := h.Peek(); hasData != false { + t.Fatal("expected Peek to return false") + } else if h.HasData() { + t.Fatal("expected HasData to return false") + } + }) +} + /* func TestHolder_Schema(t *testing.T) { t.Run("Schema", func(t *testing.T) { diff --git a/server.go b/server.go index 06c0a81e7..da6cc8a1a 100644 --- a/server.go +++ b/server.go @@ -149,9 +149,7 @@ func (s *Server) Open() error { // Don't actually load the data until after the Cluster // management starts. s.Holder.LogOutput = s.LogOutput - if err := s.Holder.Peek(); err != nil { - return fmt.Errorf("peeking at the Holder: %v", err) - } + s.Holder.Peek() // Start the BroadcastReceiver. if err := s.BroadcastReceiver.Start(s); err != nil { From 212628b2203d34fb25fd5f4bf7012ade4f0be2dc Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Fri, 10 Nov 2017 11:38:16 -0600 Subject: [PATCH 3/3] adjust tests to include Cluster.ListenForJoins() --- cluster_test.go | 11 ----------- test/cluster.go | 7 +++++++ 2 files changed, 7 insertions(+), 11 deletions(-) diff --git a/cluster_test.go b/cluster_test.go index a51823993..1e5427df0 100644 --- a/cluster_test.go +++ b/cluster_test.go @@ -308,17 +308,6 @@ func TestCluster_Resize(t *testing.T) { // TestTestCluster ensures that general cluster functionality works as expected. func TestCluster_ResizeStates(t *testing.T) { - /* test conditions: - x- single node, no data, comes up in NORMAL with topology - x- single node, in topology, comes up NORMAL - x- single node, not in topology, raises error - x- two node, no data, comes up in NORMAL, with topology - x- two node, in topology, comes up NORMAL - x- two node, STARTING, not in topology, raises error - x- two node, NORMAL, not in topology, triggers resize - x- resize of nodes with data moves data appropriately - */ - t.Run("Single node, no data", func(t *testing.T) { tc := test.NewTestCluster(1) diff --git a/test/cluster.go b/test/cluster.go index e73289827..82c45aef6 100644 --- a/test/cluster.go +++ b/test/cluster.go @@ -255,6 +255,13 @@ func (t *TestCluster) Open() error { return err } } + + // Start the listener on the coordinator. + if len(t.Clusters) == 0 { + return nil + } + t.Clusters[0].ListenForJoins() + return nil }