From a0dda250a573fd1841d8aebe3f7255ddff911126 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Mon, 1 Oct 2018 10:11:10 -0500 Subject: [PATCH 1/6] converted to pilosa.logger --- translate.go | 16 ++++++++++++---- 1 file changed, 12 insertions(+), 4 deletions(-) diff --git a/translate.go b/translate.go index cb16de929..be640a777 100644 --- a/translate.go +++ b/translate.go @@ -69,7 +69,7 @@ type TranslateFile struct { Path string mapSize int - + logger Logger // If non-nil, data is streamed from a primary and this is a read-only store. PrimaryTranslateStore TranslateStore primaryID string // unique ID used to identify the primary store @@ -90,6 +90,12 @@ func OptTranslateFileMapSize(mapSize int) TranslateFileOption { return nil } } +func OptTranslateFileLogger(l Logger) TranslateFileOption { + return func(s *TranslateFile) error { + s.logger = l + return nil + } +} // NewTranslateFile returns a new instance of TranslateFile. func NewTranslateFile(opts ...TranslateFileOption) *TranslateFile { @@ -111,6 +117,8 @@ func NewTranslateFile(opts ...TranslateFileOption) *TranslateFile { mapSize: defaultMapSize, + logger: NopLogger, + replicationClosing: make(chan struct{}), primaryStoreEvents: make(chan primaryStoreEvent), @@ -187,12 +195,12 @@ func (s *TranslateFile) handlePrimaryStoreEvent(ev primaryStoreEvent) error { } // Stop translate store replication. - log.Printf("stop monitor replication") + s.logger.Printf("stop monitor replication") close(s.replicationClosing) s.repWG.Wait() // Set the primary node for translate store replication. - log.Printf("set primary translate store to %s", ev.id) + s.logger.Printf("set primary translate store to %s", ev.id) s.primaryID = ev.id if ev.id == "" { s.PrimaryTranslateStore = nil @@ -201,7 +209,7 @@ func (s *TranslateFile) handlePrimaryStoreEvent(ev primaryStoreEvent) error { } // Start translate store replication. Stream from primary, if available. - log.Printf("start monitor replication") + s.logger.Printf("start monitor replication") if s.PrimaryTranslateStore != nil { s.replicationClosing = make(chan struct{}) s.repWG.Add(1) From 85ebaf298d5c5f5446a62b8a447e136826611fa8 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Mon, 1 Oct 2018 11:04:47 -0500 Subject: [PATCH 2/6] removed logging from translate store replaced with error --- holder.go | 8 ++++++-- http/translator.go | 10 ++++------ http/translator_test.go | 10 +++++++--- server.go | 3 ++- server/handler_test.go | 2 +- translate.go | 13 ++++++------- 6 files changed, 26 insertions(+), 20 deletions(-) diff --git a/holder.go b/holder.go index 794104e53..80b13e3b6 100644 --- a/holder.go +++ b/holder.go @@ -52,7 +52,7 @@ type Holder struct { // Key/ID translation translateFile *TranslateFile - NewPrimaryTranslateStore func(interface{}) TranslateStore + NewPrimaryTranslateStore func(interface{}) (TranslateStore, error) // opened channel is closed once Open() completes. opened chan struct{} @@ -584,7 +584,11 @@ func (h *Holder) setPrimaryTranslateStore(node *Node) { if node != nil { nodeID = node.ID } - h.translateFile.SetPrimaryStore(nodeID, h.NewPrimaryTranslateStore(node)) + ts, err := h.NewPrimaryTranslateStore(node) + if err != nil { + h.Logger.Printf("setPrimaryTranslateStore: %s", err) + } + h.translateFile.SetPrimaryStore(nodeID, ts) } // holderSyncer is an active anti-entropy tool that compares the local holder diff --git a/http/translator.go b/http/translator.go index 7ca85a33f..f5e9f3ed3 100644 --- a/http/translator.go +++ b/http/translator.go @@ -6,7 +6,6 @@ import ( "fmt" "io" "io/ioutil" - "log" "net/http" "net/url" "strconv" @@ -27,13 +26,12 @@ type translateStore struct { // NewTranslateStore returns a new instance of TranslateStore based on node. // DEPRECATED: Providing a string url to this function is being deprecated. Instead, // provide a *pilosa.Node. -func NewTranslateStore(node interface{}) pilosa.TranslateStore { +func NewTranslateStore(node interface{}) (pilosa.TranslateStore, error) { var n *pilosa.Node switch v := node.(type) { case string: - log.Printf("WARNING: providing a string url to NewTranslateStore() has been deprecated.") if uri, err := pilosa.NewURIFromAddress(v); err != nil { - log.Println(errors.Wrap(err, "creating uri")) + return nil, errors.Wrap(err, "creating uri") } else { n = &pilosa.Node{ ID: v, @@ -43,9 +41,9 @@ func NewTranslateStore(node interface{}) pilosa.TranslateStore { case *pilosa.Node: n = v default: - log.Printf("WARNING: a *pilosa.Node is the only type supported by NewTranslateStore().") + return nil, errors.New("*pilosa.Node is the only type supported by NewTranslateStore().") } - return &translateStore{node: n} + return &translateStore{node: n}, nil } // TranslateColumnsToUint64 is not currently implemented. diff --git a/http/translator_test.go b/http/translator_test.go index 18ba7c998..496d85cc9 100644 --- a/http/translator_test.go +++ b/http/translator_test.go @@ -42,7 +42,7 @@ func TestTranslateStore_Reader(t *testing.T) { } // Connect to server and stream all available data. - store := http.NewTranslateStore(primary.URL()) + store, _ := http.NewTranslateStore(primary.URL()) // Wait to ensure writes make it to translate store time.Sleep(500 * time.Millisecond) @@ -96,7 +96,7 @@ func TestTranslateStore_Reader(t *testing.T) { // Connect to server and begin streaming. ctx, cancel := context.WithCancel(context.Background()) - store := http.NewTranslateStore(primary.URL()) + store, _ := http.NewTranslateStore(primary.URL()) if _, err := store.Reader(ctx, 0); err != nil { t.Fatal(err) } @@ -124,7 +124,11 @@ func TestTranslateStore_Reader(t *testing.T) { primary := test.MustRunCluster(t, 1, []server.CommandOption{opts})[0] defer primary.Close() - _, err := http.NewTranslateStore(primary.URL()).Reader(context.Background(), 0) + ts, err := http.NewTranslateStore(primary.URL()) + if err != nil { + t.Fatalf("unexpected error: %s", err) + } + _, err = ts.Reader(context.Background(), 0) if err != pilosa.ErrNotImplemented { t.Fatalf("unexpected error: %s", err) } diff --git a/server.go b/server.go index 61e4838a1..ce817ba35 100644 --- a/server.go +++ b/server.go @@ -168,7 +168,8 @@ func OptServerPrimaryTranslateStore(store TranslateStore) ServerOption { } } -func OptServerPrimaryTranslateStoreFunc(tf func(interface{}) TranslateStore) ServerOption { +func OptServerPrimaryTranslateStoreFunc(tf func(interface{}) (TranslateStore, error)) ServerOption { + return func(s *Server) error { s.holder.NewPrimaryTranslateStore = tf return nil diff --git a/server/handler_test.go b/server/handler_test.go index 2c1aa9868..c0144cb6c 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -680,7 +680,7 @@ func TestClusterTranslator(t *testing.T) { cluster[0] = test.NewCommandNode(true) cluster[0].Config.Gossip.Port = "0" cluster[0].Start() - httpTranslateStore := http.NewTranslateStore(cluster[0].URL()) + httpTranslateStore, _ := http.NewTranslateStore(cluster[0].URL()) cluster[1] = test.NewCommandNode(false, server.OptCommandServerOptions( pilosa.OptServerPrimaryTranslateStore(httpTranslateStore), diff --git a/translate.go b/translate.go index be640a777..b875c415d 100644 --- a/translate.go +++ b/translate.go @@ -8,7 +8,6 @@ import ( "fmt" "io" "io/ioutil" - "log" "os" "path/filepath" "sync" @@ -371,14 +370,14 @@ func (s *TranslateFile) monitorReplication() { // Keep attempting to replicate until the store closes. for { if err := s.replicate(ctx); err != nil { - log.Printf("pilosa: replication error: %s", err) + s.logger.Printf("pilosa: replication error: %s", err) } select { case <-ctx.Done(): return case <-time.After(s.replicationRetryInterval): - log.Printf("pilosa: reconnecting to primary replica") + s.logger.Printf("pilosa: reconnecting to primary replica") } } } @@ -386,7 +385,7 @@ func (s *TranslateFile) monitorReplication() { // monitorPrimaryStoreEvents is executed in a separate goroutine and listens for changes // to the primary store assignment. func (s *TranslateFile) monitorPrimaryStoreEvents() { - log.Printf("monitor primary store events") + s.logger.Printf("monitor primary store events") // Keep handling events until the store closes. for { select { @@ -394,7 +393,7 @@ func (s *TranslateFile) monitorPrimaryStoreEvents() { return case ev := <-s.primaryStoreEvents: if err := s.handlePrimaryStoreEvent(ev); err != nil { - log.Printf("handle primary store event") + s.logger.Printf("handle primary store event") } } } @@ -404,7 +403,7 @@ func (s *TranslateFile) replicate(ctx context.Context) error { off := s.size() // Connect to remote primary. - log.Printf("pilosa: replicating from offset %d", off) + s.logger.Printf("pilosa: replicating from offset %d", off) rc, err := s.PrimaryTranslateStore.Reader(ctx, off) if err != nil { return err @@ -1119,7 +1118,7 @@ var nopTStore TranslateStore = nopTranslateStore{} // newNopTranslateStore returns a translate store which does nothing. It returns a global // object to avoid unnecessary allocations. -func newNopTranslateStore(interface{}) TranslateStore { return nopTStore } +func newNopTranslateStore(interface{}) (TranslateStore, error) { return nopTStore, nil } // nopTranslateStore represents a no-op implementation of the TranslateStore interface. type nopTranslateStore struct{} From 0df193cf98f0d2b022c47a01423488d0525ac1e3 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Wed, 17 Oct 2018 10:29:58 -0500 Subject: [PATCH 3/6] revert to existing api with panic per jaffee --- holder.go | 7 ++----- http/translator.go | 9 ++++----- http/translator_test.go | 11 ++++------- server.go | 6 ++++-- server/handler_test.go | 2 +- translate.go | 2 +- 6 files changed, 16 insertions(+), 21 deletions(-) diff --git a/holder.go b/holder.go index 899c25944..0a1af820b 100644 --- a/holder.go +++ b/holder.go @@ -52,7 +52,7 @@ type Holder struct { // Key/ID translation translateFile *TranslateFile - NewPrimaryTranslateStore func(interface{}) (TranslateStore, error) + NewPrimaryTranslateStore func(interface{}) TranslateStore // opened channel is closed once Open() completes. opened chan struct{} @@ -590,10 +590,7 @@ func (h *Holder) setPrimaryTranslateStore(node *Node) { if node != nil { nodeID = node.ID } - ts, err := h.NewPrimaryTranslateStore(node) - if err != nil { - h.Logger.Printf("setPrimaryTranslateStore: %s", err) - } + ts := h.NewPrimaryTranslateStore(node) h.translateFile.SetPrimaryStore(nodeID, ts) } diff --git a/http/translator.go b/http/translator.go index f5e9f3ed3..f83bde113 100644 --- a/http/translator.go +++ b/http/translator.go @@ -11,7 +11,6 @@ import ( "strconv" "github.com/pilosa/pilosa" - "github.com/pkg/errors" ) // Ensure implementation implements inteface. @@ -26,12 +25,12 @@ type translateStore struct { // NewTranslateStore returns a new instance of TranslateStore based on node. // DEPRECATED: Providing a string url to this function is being deprecated. Instead, // provide a *pilosa.Node. -func NewTranslateStore(node interface{}) (pilosa.TranslateStore, error) { +func NewTranslateStore(node interface{}) pilosa.TranslateStore { var n *pilosa.Node switch v := node.(type) { case string: if uri, err := pilosa.NewURIFromAddress(v); err != nil { - return nil, errors.Wrap(err, "creating uri") + panic("bad uri for translatestore in deprecated api") } else { n = &pilosa.Node{ ID: v, @@ -41,9 +40,9 @@ func NewTranslateStore(node interface{}) (pilosa.TranslateStore, error) { case *pilosa.Node: n = v default: - return nil, errors.New("*pilosa.Node is the only type supported by NewTranslateStore().") + panic("*pilosa.Node is the only type supported by NewTranslateStore().") } - return &translateStore{node: n}, nil + return &translateStore{node: n} } // TranslateColumnsToUint64 is not currently implemented. diff --git a/http/translator_test.go b/http/translator_test.go index 496d85cc9..ce96b3c2c 100644 --- a/http/translator_test.go +++ b/http/translator_test.go @@ -42,7 +42,7 @@ func TestTranslateStore_Reader(t *testing.T) { } // Connect to server and stream all available data. - store, _ := http.NewTranslateStore(primary.URL()) + store := http.NewTranslateStore(primary.URL()) // Wait to ensure writes make it to translate store time.Sleep(500 * time.Millisecond) @@ -96,7 +96,7 @@ func TestTranslateStore_Reader(t *testing.T) { // Connect to server and begin streaming. ctx, cancel := context.WithCancel(context.Background()) - store, _ := http.NewTranslateStore(primary.URL()) + store := http.NewTranslateStore(primary.URL()) if _, err := store.Reader(ctx, 0); err != nil { t.Fatal(err) } @@ -124,11 +124,8 @@ func TestTranslateStore_Reader(t *testing.T) { primary := test.MustRunCluster(t, 1, []server.CommandOption{opts})[0] defer primary.Close() - ts, err := http.NewTranslateStore(primary.URL()) - if err != nil { - t.Fatalf("unexpected error: %s", err) - } - _, err = ts.Reader(context.Background(), 0) + ts := http.NewTranslateStore(primary.URL()) + _, err := ts.Reader(context.Background(), 0) if err != pilosa.ErrNotImplemented { t.Fatalf("unexpected error: %s", err) } diff --git a/server.go b/server.go index 77593292c..3815f1bf4 100644 --- a/server.go +++ b/server.go @@ -168,7 +168,7 @@ func OptServerPrimaryTranslateStore(store TranslateStore) ServerOption { } } -func OptServerPrimaryTranslateStoreFunc(tf func(interface{}) (TranslateStore, error)) ServerOption { +func OptServerPrimaryTranslateStoreFunc(tf func(interface{}) TranslateStore) ServerOption { return func(s *Server) error { s.holder.NewPrimaryTranslateStore = tf @@ -237,7 +237,9 @@ func OptServerClusterHasher(h Hasher) ServerOption { func OptServerTranslateFileMapSize(mapSize int) ServerOption { return func(s *Server) error { - s.holder.translateFile = NewTranslateFile(OptTranslateFileMapSize(mapSize)) + s.holder.translateFile = NewTranslateFile( + OptTranslateFileMapSize(mapSize), + OptTranslateFileLogger(s.logger)) return nil } } diff --git a/server/handler_test.go b/server/handler_test.go index aac2447f1..9e76d77b9 100644 --- a/server/handler_test.go +++ b/server/handler_test.go @@ -696,7 +696,7 @@ func TestClusterTranslator(t *testing.T) { cluster[0] = test.NewCommandNode(true) cluster[0].Config.Gossip.Port = "0" cluster[0].Start() - httpTranslateStore, _ := http.NewTranslateStore(cluster[0].URL()) + httpTranslateStore := http.NewTranslateStore(cluster[0].URL()) cluster[1] = test.NewCommandNode(false, server.OptCommandServerOptions( pilosa.OptServerPrimaryTranslateStore(httpTranslateStore), diff --git a/translate.go b/translate.go index b875c415d..86ed6c914 100644 --- a/translate.go +++ b/translate.go @@ -1118,7 +1118,7 @@ var nopTStore TranslateStore = nopTranslateStore{} // newNopTranslateStore returns a translate store which does nothing. It returns a global // object to avoid unnecessary allocations. -func newNopTranslateStore(interface{}) (TranslateStore, error) { return nopTStore, nil } +func newNopTranslateStore(interface{}) TranslateStore { return nopTStore } // nopTranslateStore represents a no-op implementation of the TranslateStore interface. type nopTranslateStore struct{} From e5e2237fc0e7a51baf3face1fa30686e933c72fb Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Wed, 17 Oct 2018 11:12:37 -0500 Subject: [PATCH 4/6] changed order of logger options --- server.go | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/server.go b/server.go index 3815f1bf4..db8f38895 100644 --- a/server.go +++ b/server.go @@ -237,9 +237,7 @@ func OptServerClusterHasher(h Hasher) ServerOption { func OptServerTranslateFileMapSize(mapSize int) ServerOption { return func(s *Server) error { - s.holder.translateFile = NewTranslateFile( - OptTranslateFileMapSize(mapSize), - OptTranslateFileLogger(s.logger)) + s.holder.translateFile = NewTranslateFile( OptTranslateFileMapSize(mapSize)) return nil } } @@ -273,6 +271,7 @@ func NewServer(opts ...ServerOption) (*Server, error) { return nil, errors.Wrap(err, "applying option") } } + s.holder.translateFile.logger = s.logger path, err := expandDirName(s.dataDir) if err != nil { From c0a3b7a781d279e24ed24979ef56656344b2fd2a Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Wed, 17 Oct 2018 12:35:00 -0500 Subject: [PATCH 5/6] fmt --- server.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/server.go b/server.go index 2e566b919..a4cb8fc8b 100644 --- a/server.go +++ b/server.go @@ -237,7 +237,7 @@ func OptServerClusterHasher(h Hasher) ServerOption { func OptServerTranslateFileMapSize(mapSize int) ServerOption { return func(s *Server) error { - s.holder.translateFile = NewTranslateFile( OptTranslateFileMapSize(mapSize)) + s.holder.translateFile = NewTranslateFile(OptTranslateFileMapSize(mapSize)) return nil } } From 0d189a47a2cc022c641f729152ff87fb79f534e4 Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Wed, 17 Oct 2018 17:58:52 -0500 Subject: [PATCH 6/6] added v2 notation --- http/translator.go | 1 + 1 file changed, 1 insertion(+) diff --git a/http/translator.go b/http/translator.go index f83bde113..3a2acae04 100644 --- a/http/translator.go +++ b/http/translator.go @@ -25,6 +25,7 @@ type translateStore struct { // NewTranslateStore returns a new instance of TranslateStore based on node. // DEPRECATED: Providing a string url to this function is being deprecated. Instead, // provide a *pilosa.Node. +// TODO (2.0) Refactor to avoid panic func NewTranslateStore(node interface{}) pilosa.TranslateStore { var n *pilosa.Node switch v := node.(type) {