From 28acc29a109433c69927889624099c2aed740b67 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Mon, 23 Apr 2018 10:37:20 -0500 Subject: [PATCH 1/5] refactoring pilosa/server trying to separate internal an external concerns in pilosa.Server - it should handle Cluster, Holder, etc. while pilosa/server handles things with external deps - e.g. Logger, Stats, Handler, etc. Using functional options in pilosa.Server now. --- broadcast_test.go | 11 +- client_test.go | 3 +- cmd/server.go | 54 +--------- cmd/server_test.go | 12 +-- ctl/common.go | 2 +- ctl/server.go | 4 - handler.go | 7 -- holder.go | 1 - holder_test.go | 5 +- server.go | 228 +++++++++++++++++++++++++---------------- server/cluster_test.go | 69 +------------ server/server.go | 215 ++++++++++++++++++++++++-------------- server/server_test.go | 6 +- test/executor.go | 3 +- test/pilosa.go | 23 ++--- test/test.go | 7 -- 16 files changed, 311 insertions(+), 339 deletions(-) diff --git a/broadcast_test.go b/broadcast_test.go index f4ff4c0fa..187e50b4e 100644 --- a/broadcast_test.go +++ b/broadcast_test.go @@ -15,12 +15,16 @@ package pilosa_test import ( + "bytes" "reflect" "testing" + "io/ioutil" + "github.com/gogo/protobuf/proto" "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/internal" + "github.com/pilosa/pilosa/server" ) // Ensure a message can be marshaled and unmarshaled. @@ -52,11 +56,12 @@ func testMessageMarshal(t *testing.T, m proto.Message) { // Ensure that BroadcastReceiver can register a BroadcastHandler. func TestBroadcast_BroadcastReceiver(t *testing.T) { - - s, err := pilosa.NewServer() + com := server.NewCommand(bytes.NewBuffer([]byte{}), ioutil.Discard, ioutil.Discard) + err := com.SetupServer() // this test shouldn't need to import pilosa/server just to set up the Server, but it really shouldn't need to setup the Server at all. The Server should not be the implementation of Broadcast* TODO if err != nil { - t.Fatalf("getting new server: %v", err) + t.Fatalf("setting up server: %v", err) } + s := com.Server sbr := NewSimpleBroadcastReceiver() sbh := NewSimpleBroadcastHandler() diff --git a/client_test.go b/client_test.go index 4bd0bade6..d11ab16a9 100644 --- a/client_test.go +++ b/client_test.go @@ -26,6 +26,7 @@ import ( "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/internal" "github.com/pilosa/pilosa/pql" + "github.com/pilosa/pilosa/server" "github.com/pilosa/pilosa/test" ) @@ -47,7 +48,7 @@ func createCluster(c *pilosa.Cluster) ([]*test.Server, []*test.Holder) { var defaultClient *http.Client func init() { - defaultClient = pilosa.GetHTTPClient(nil) + defaultClient = server.GetHTTPClient(nil) } diff --git a/cmd/server.go b/cmd/server.go index 81633ea01..93ecb543b 100644 --- a/cmd/server.go +++ b/cmd/server.go @@ -15,17 +15,11 @@ package cmd import ( - "fmt" "io" - "os" - "os/signal" - "runtime/pprof" - "syscall" - "time" + "github.com/pkg/errors" "github.com/spf13/cobra" - "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/ctl" "github.com/pilosa/pilosa/server" ) @@ -45,52 +39,10 @@ It will load existing data from the configured directory and start listening for client connections on the configured port.`, RunE: func(cmd *cobra.Command, args []string) error { - // Set up the logger. - if err := Server.SetupLogger(); err != nil { - return fmt.Errorf("error setting up the logger: %v", err) - } - logger := Server.Server.Logger - logger.Printf("Pilosa %s, build time %s\n", pilosa.Version, pilosa.BuildTime) - - // Start CPU profiling. - if Server.CPUProfile != "" { - f, err := os.Create(Server.CPUProfile) - if err != nil { - return fmt.Errorf("create cpu profile: %v", err) - } - defer f.Close() - - fmt.Fprintln(Server.Stderr, "Starting cpu profile") - pprof.StartCPUProfile(f) - time.AfterFunc(Server.CPUTime, func() { - fmt.Fprintln(Server.Stderr, "Stopping cpu profile") - pprof.StopCPUProfile() - f.Close() - }) - } - - // Execute the program. if err := Server.Run(); err != nil { - return fmt.Errorf("error running server: %v", err) + return errors.Wrap(err, "running server") } - - // First SIGKILL causes server to shut down gracefully. - c := make(chan os.Signal, 2) - signal.Notify(c, os.Interrupt, syscall.SIGTERM) - select { - case sig := <-c: - logger.Printf("Received %s; gracefully shutting down...\n", sig.String()) - - // Second signal causes a hard shutdown. - go func() { <-c; os.Exit(1) }() - - if err := Server.Close(); err != nil { - return err - } - case <-Server.Done: - logger.Printf("Server closed externally") - } - return nil + return errors.Wrap(Server.Wait(), "waiting on Server") }, } diff --git a/cmd/server_test.go b/cmd/server_test.go index e69ce1134..abbe8d7a4 100644 --- a/cmd/server_test.go +++ b/cmd/server_test.go @@ -37,8 +37,6 @@ func TestServerHelp(t *testing.T) { func TestServerConfig(t *testing.T) { actualDataDir, err := ioutil.TempDir("", "") failErr(t, err, "making data dir") - profFile, err := ioutil.TempFile("", "") - failErr(t, err, "making temp file") logFile, err := ioutil.TempFile("", "") failErr(t, err, "making log file") tests := []commandTest{ @@ -93,7 +91,7 @@ func TestServerConfig(t *testing.T) { // TEST 2 { args: []string{"server", "--log-path", logFile.Name(), "--cluster.disabled", "true"}, - env: map[string]string{"PILOSA_PROFILE_CPU_TIME": "1m"}, + env: map[string]string{}, cfgFileContent: ` bind = "localhost:19444" data-dir = "` + actualDataDir + `" @@ -103,9 +101,6 @@ func TestServerConfig(t *testing.T) { ] [anti-entropy] interval = "11m0s" - [profile] - cpu = "` + profFile.Name() + `" - cpu-time = "35s" [metric] service = "statsd" host = "127.0.0.1:8125" @@ -114,8 +109,6 @@ func TestServerConfig(t *testing.T) { v := validator{} v.Check(cmd.Server.Config.Cluster.Hosts, []string{"localhost:19444"}) v.Check(cmd.Server.Config.AntiEntropy.Interval, toml.Duration(time.Minute*11)) - v.Check(cmd.Server.CPUProfile, profFile.Name()) - v.Check(cmd.Server.CPUTime, time.Minute) v.Check(cmd.Server.Config.LogPath, logFile.Name()) v.Check(cmd.Server.Config.Metric.Service, "statsd") v.Check(cmd.Server.Config.Metric.Host, "127.0.0.1:8125") @@ -147,6 +140,9 @@ func TestServerConfig(t *testing.T) { case <-cmd.Server.Started: case <-executed: } + if execErr != nil { + t.Fatalf("executing server command: %v", execErr) + } err := cmd.Server.Close() failErr(t, err, "closing pilosa server command") <-executed diff --git a/ctl/common.go b/ctl/common.go index 6ebcd3132..3629e233e 100644 --- a/ctl/common.go +++ b/ctl/common.go @@ -49,7 +49,7 @@ func CommandClient(cmd CommandWithTLSSupport) (*pilosa.InternalHTTPClient, error InsecureSkipVerify: tlsConfig.SkipVerify, } } - client, err := pilosa.NewInternalHTTPClient(cmd.TLSHost(), pilosa.GetHTTPClient(TLSConfig)) + client, err := pilosa.NewInternalHTTPClient(cmd.TLSHost(), server.GetHTTPClient(TLSConfig)) if err != nil { return nil, err } diff --git a/ctl/server.go b/ctl/server.go index 1c60a7c9c..4816f6170 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -61,8 +61,4 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) { flags.StringVarP(&srv.Config.Metric.Host, "metric.host", "", srv.Config.Metric.Host, "Default URI to send metrics.") flags.DurationVarP((*time.Duration)(&srv.Config.Metric.PollInterval), "metric.poll-interval", "", (time.Duration)(srv.Config.Metric.PollInterval), "Polling interval metrics.") flags.BoolVarP((&srv.Config.Metric.Diagnostics), "metric.diagnostics", "", srv.Config.Metric.Diagnostics, "Enabled diagnostics reporting.") - - // CPU Profiling - flags.StringVarP(&srv.CPUProfile, "profile.cpu", "", "", "Where to store CPU profile.") - flags.DurationVarP(&srv.CPUTime, "profile.cpu-time", "", 30*time.Second, "CPU profile duration.") } diff --git a/handler.go b/handler.go index bd97c1d1a..52cbff094 100644 --- a/handler.go +++ b/handler.go @@ -15,7 +15,6 @@ package pilosa import ( - "context" "encoding/json" "expvar" "fmt" @@ -35,7 +34,6 @@ import ( "github.com/gogo/protobuf/proto" "github.com/gorilla/mux" "github.com/pilosa/pilosa/internal" - "github.com/pilosa/pilosa/pql" "github.com/pkg/errors" ) @@ -45,11 +43,6 @@ type Handler struct { FileSystem FileSystem - // The execution engine for running queries. - Executor interface { - Execute(context context.Context, index string, query *pql.Query, slices []uint64, opt *ExecOptions) ([]interface{}, error) - } - Logger Logger // Keeps the query argument validators for each handler diff --git a/holder.go b/holder.go index fa93767a3..2279a3e19 100644 --- a/holder.go +++ b/holder.go @@ -530,7 +530,6 @@ func (h *Holder) setFileLimit() { func (h *Holder) loadNodeID() (string, error) { idPath := path.Join(h.Path, "ID") nodeID := "" - h.Logger.Printf("load NodeID: %s", idPath) if err := os.MkdirAll(h.Path, 0777); err != nil { return "", err diff --git a/holder_test.go b/holder_test.go index 8b1f6166a..7053fa1f2 100644 --- a/holder_test.go +++ b/holder_test.go @@ -25,6 +25,7 @@ import ( "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/pql" + "github.com/pilosa/pilosa/server" "github.com/pilosa/pilosa/test" ) @@ -367,7 +368,7 @@ func TestHolder_DeleteIndex(t *testing.T) { // Ensure holder can sync with a remote holder. func TestHolderSyncer_SyncHolder(t *testing.T) { cluster := test.NewCluster(2) - client := pilosa.GetHTTPClient(nil) + client := server.GetHTTPClient(nil) // Create a local holder. hldr0 := test.MustOpenHolder() defer hldr0.Close() @@ -451,7 +452,7 @@ func TestHolderSyncer_SyncHolder(t *testing.T) { Holder: hldr0.Holder, Node: cluster.Nodes[0], Cluster: cluster, - RemoteClient: pilosa.GetHTTPClient(nil), + RemoteClient: server.GetHTTPClient(nil), Stats: pilosa.NopStatsClient, } diff --git a/server.go b/server.go index 6d25e3320..f022d562a 100644 --- a/server.go +++ b/server.go @@ -55,55 +55,154 @@ func OptServerLogger(l Logger) ServerOption { } } +func OptServerReplicaN(n int) ServerOption { + return func(s *Server) error { + s.Cluster.ReplicaN = n + return nil + } +} + +func OptServerDataDir(dir string) ServerOption { + return func(s *Server) error { + s.Cluster.Path = dir + s.Holder.Path = dir + return nil + } +} + +func OptServerAttrStoreFunc(af func(string) AttrStore) ServerOption { + return func(s *Server) error { + s.NewAttrStore = af + s.Holder.NewAttrStore = af + return nil + } +} + +func OptServerAntiEntropyInterval(interval time.Duration) ServerOption { + return func(s *Server) error { + s.AntiEntropyInterval = interval + return nil + } +} + +func OptServerLongQueryTime(dur time.Duration) ServerOption { + return func(s *Server) error { + s.Cluster.LongQueryTime = dur + return nil + } +} + +func OptServerHandler(h *Handler) ServerOption { + return func(s *Server) error { + s.Handler = h + return nil + } +} + +func OptServerMaxWritesPerRequest(n int) ServerOption { + return func(s *Server) error { + s.MaxWritesPerRequest = n + return nil + } +} + +func OptServerMetricInterval(dur time.Duration) ServerOption { + return func(s *Server) error { + s.MetricInterval = dur + return nil + } +} + +func OptServerSystemInfo(si SystemInfo) ServerOption { + return func(s *Server) error { + s.SystemInfo = si + return nil + } +} + +func OptServerGCNotifier(gcn GCNotifier) ServerOption { + return func(s *Server) error { + s.GCNotifier = gcn + return nil + } +} + +func OptServerRemoteClient(c *http.Client) ServerOption { + return func(s *Server) error { + s.RemoteClient = c + s.Cluster.RemoteClient = c + return nil + } +} + +func OptServerStatsClient(sc StatsClient) ServerOption { + return func(s *Server) error { + s.Holder.Stats = sc + return nil + } +} + +func OptServerDiagnosticsInterval(dur time.Duration) ServerOption { + return func(s *Server) error { + s.DiagnosticInterval = dur + return nil + } +} + +func OptServerListener(ln net.Listener) ServerOption { + return func(s *Server) error { + s.ln = ln + + return nil + } +} + +func OptServerURI(uri *URI) ServerOption { + return func(s *Server) error { + s.URI = *uri + return nil + } +} + // Server represents a holder wrapped by a running HTTP server. type Server struct { - ln net.Listener - // Close management. wg sync.WaitGroup closing chan struct{} - // Data storage and HTTP interface. - Holder *Holder + // Internal + Holder *Holder + Cluster *Cluster + diagnostics *DiagnosticsCollector + + // External Handler *Handler Broadcaster Broadcaster BroadcastReceiver BroadcastReceiver Gossiper Gossiper RemoteClient *http.Client + SystemInfo SystemInfo + GCNotifier GCNotifier + NewAttrStore func(string) AttrStore + Logger Logger + TLS *tls.Config + ln net.Listener - // Cluster configuration. - Network string - NodeID string - URI URI - Cluster *Cluster - diagnostics *DiagnosticsCollector - SystemInfo SystemInfo - - GCNotifier GCNotifier - - NewAttrStore func(string) AttrStore - - // Background monitoring intervals. + NodeID string + URI URI AntiEntropyInterval time.Duration MetricInterval time.Duration DiagnosticInterval time.Duration - - // TLS configuration - TLS *tls.Config - - // Misc options. MaxWritesPerRequest int - Logger Logger - defaultClient InternalClient } // NewServer returns a new instance of Server. func NewServer(opts ...ServerOption) (*Server, error) { s := &Server{ - closing: make(chan struct{}), - + closing: make(chan struct{}), + Cluster: NewCluster(), Holder: NewHolder(), Handler: NewHandler(), Broadcaster: NopBroadcaster, @@ -111,8 +210,6 @@ func NewServer(opts ...ServerOption) (*Server, error) { diagnostics: NewDiagnosticsCollector(DefaultDiagnosticServer), SystemInfo: NewNopSystemInfo(), - Network: "tcp", - GCNotifier: NopGCNotifier, NewAttrStore: NewNopAttrStore, @@ -131,20 +228,25 @@ func NewServer(opts ...ServerOption) (*Server, error) { } } - s.Handler.API = NewAPI() - s.Handler.API.Holder = s.Holder + s.Holder.Logger = s.Logger + s.Holder.Stats.SetLogger(s.Logger) + + s.Cluster.Logger = s.Logger + s.Cluster.Holder = s.Holder + s.Cluster.RemoteClient = s.RemoteClient + // update URI port with actual listener port. TODO this should probably be done outside of here. + if s.URI.Port() == 0 { + s.URI.SetPort(uint16(s.ln.Addr().(*net.TCPAddr).Port)) + } return s, nil } // Open opens and initializes the server. func (s *Server) Open() error { - s.Handler.API.Logger = s.Logger // TODO do this in NewServer with functional options s.Logger.Printf("open server") // s.ln can be configured prior to Open() via s.OpenListener(). if s.ln == nil { - if err := s.OpenListener(); err != nil { - return err - } + return errors.New("Must pass a listener option to NewServer") } // Get or create NodeID. @@ -181,13 +283,13 @@ func (s *Server) Open() error { s.Cluster.MaxWritesPerRequest = s.MaxWritesPerRequest // Initialize HTTP handler. + s.Handler.API.Holder = s.Holder s.Handler.API.Broadcaster = s.Broadcaster s.Handler.API.BroadcastHandler = s s.Handler.API.StatusHandler = s s.Handler.API.URI = s.URI s.Handler.API.Cluster = s.Cluster s.Handler.API.Executor = e - s.Handler.Executor = e // Initialize Holder. s.Holder.Broadcaster = s.Broadcaster @@ -234,43 +336,6 @@ func (s *Server) Open() error { return nil } -// OpenListener opens a listener for the Server. -func (s *Server) OpenListener() error { - s.Logger.Printf("open server listener: %s", s.URI) - if s.ln != nil { - return fmt.Errorf("a listener already exists for server: %s", s.URI) - } - - var ln net.Listener - var err error - - // If bind URI has the https scheme, enable TLS - if s.URI.Scheme() == "https" && s.TLS != nil { - ln, err = tls.Listen("tcp", s.URI.HostPort(), s.TLS) - if err != nil { - return err - } - } else if s.URI.Scheme() == "http" { - // Open HTTP listener to determine port (if specified as :0). - ln, err = net.Listen(s.Network, s.URI.HostPort()) - if err != nil { - return fmt.Errorf("net.Listen: %v", err) - } - } else { - return fmt.Errorf("unsupported scheme: %s", s.URI.Scheme()) - } - - s.ln = ln - - if s.URI.Port() == 0 { - // If the port is 0, it is set automatically. - // Find out automatically set port and update the host. - s.URI.SetPort(uint16(s.ln.Addr().(*net.TCPAddr).Port)) - } - - return nil -} - // Close closes the server and waits for it to shutdown. func (s *Server) Close() error { // Notify goroutines to stop. @@ -311,25 +376,6 @@ func (s *Server) Addr() net.Addr { } return s.ln.Addr() } -func GetHTTPClient(t *tls.Config) *http.Client { - transport := &http.Transport{ - Proxy: http.ProxyFromEnvironment, - DialContext: (&net.Dialer{ - Timeout: 30 * time.Second, - KeepAlive: 30 * time.Second, - DualStack: true, - }).DialContext, - MaxIdleConns: 1000, - MaxIdleConnsPerHost: 200, - IdleConnTimeout: 90 * time.Second, - TLSHandshakeTimeout: 10 * time.Second, - ExpectContinueTimeout: 1 * time.Second, - } - if t != nil { - transport.TLSClientConfig = t - } - return &http.Client{Transport: transport} -} func (s *Server) monitorAntiEntropy() { ticker := time.NewTicker(s.AntiEntropyInterval) diff --git a/server/cluster_test.go b/server/cluster_test.go index 5b0321285..1a64ce305 100644 --- a/server/cluster_test.go +++ b/server/cluster_test.go @@ -26,81 +26,16 @@ import ( "golang.org/x/sync/errgroup" "github.com/pilosa/pilosa" - "github.com/pilosa/pilosa/gossip" "github.com/pilosa/pilosa/test" ) // Ensure program can send/receive broadcast messages. func TestMain_SendReceiveMessage(t *testing.T) { - - m0 := test.MustRunMain() + ms := test.MustRunMainWithCluster(t, 2) + m0, m1 := ms[0], ms[1] defer m0.Close() - - m1 := test.MustRunMain() defer m1.Close() - // Update cluster config - m0.Server.Cluster.Nodes = []*pilosa.Node{ - {ID: m0.Server.NodeID, URI: m0.Server.URI}, - {ID: m1.Server.NodeID, URI: m1.Server.URI}, - } - m1.Server.Cluster.Nodes = m0.Server.Cluster.Nodes - - // Configure node0 - - // get the host portion of addr to use for binding - m0.Config.Gossip.Port = "0" - m0.Config.Gossip.Seeds = []string{} - - m0.Server.Cluster.Coordinator = m0.Server.NodeID - m0.Server.Cluster.Topology = &pilosa.Topology{NodeIDs: []string{m0.Server.NodeID, m1.Server.NodeID}} - m0.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m0.Server.Logger) - gossipMemberSet0, err := gossip.NewGossipMemberSet(m0.Server.URI.HostPort(), m0.Server.URI.Host(), m0.Config.Gossip, m0.Server) - if err != nil { - t.Fatal(err) - } - m0.Server.Cluster.MemberSet = gossipMemberSet0 - m0.Server.Broadcaster = m0.Server - m0.Server.Gossiper = gossipMemberSet0 - m0.Server.Handler.API.Broadcaster = m0.Server.Broadcaster - m0.Server.Holder.Broadcaster = m0.Server.Broadcaster - m0.Server.BroadcastReceiver = gossipMemberSet0 - - if err := m0.Server.BroadcastReceiver.Start(m0.Server); err != nil { - t.Fatal(err) - } - // Open Cluster management. - if err := m0.Server.Cluster.Open(); err != nil { - t.Fatal(err) - } - - // Configure node1 - - // get the host portion of addr to use for binding - m1.Config.Gossip.Port = "0" - m1.Config.Gossip.Seeds = []string{gossipMemberSet0.GetBindAddr()} - - m1.Server.Cluster.Coordinator = m0.Server.NodeID - m1.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m1.Server.Logger) - gossipMemberSet1, err := gossip.NewGossipMemberSet(m1.Server.URI.HostPort(), m1.Server.URI.Host(), m1.Config.Gossip, m1.Server) - if err != nil { - t.Fatal(err) - } - m1.Server.Cluster.MemberSet = gossipMemberSet1 - m1.Server.Broadcaster = m1.Server - m1.Server.Gossiper = gossipMemberSet1 - m1.Server.Handler.API.Broadcaster = m1.Server.Broadcaster - m1.Server.Holder.Broadcaster = m1.Server.Broadcaster - m1.Server.BroadcastReceiver = gossipMemberSet1 - - if err := m1.Server.BroadcastReceiver.Start(m1.Server); err != nil { - t.Fatal(err) - } - // Open Cluster management. - if err := m1.Server.Cluster.Open(); err != nil { - t.Fatal(err) - } - m0.Server.Cluster.SetState(pilosa.ClusterStateNormal) m1.Server.Cluster.SetState(pilosa.ClusterStateNormal) diff --git a/server/server.go b/server/server.go index 2606ff248..829676fd5 100644 --- a/server/server.go +++ b/server/server.go @@ -24,10 +24,14 @@ import ( "io" "log" "math/rand" + "net" + "net/http" "os" + "os/signal" "path/filepath" "strconv" "strings" + "syscall" "time" "crypto/tls" @@ -46,10 +50,10 @@ func init() { rand.Seed(time.Now().UTC().UnixNano()) } -const ( - // DefaultDataDir is the default data directory. - DefaultDataDir = "~/.pilosa" -) +type loggerLogger interface { + pilosa.Logger + Logger() *log.Logger +} // Command represents the state of the pilosa server command. type Command struct { @@ -58,10 +62,6 @@ type Command struct { // Configuration. Config *Config - // Profiling options. - CPUProfile string - CPUTime time.Duration - // Gossip transport GossipTransport *gossip.Transport @@ -75,18 +75,12 @@ type Command struct { // Passed to the Gossip implementation. logOutput io.Writer - logger *log.Logger + logger loggerLogger } // NewCommand returns a new instance of Main. func NewCommand(stdin io.Reader, stdout, stderr io.Writer) *Command { - s, err := pilosa.NewServer() - if err != nil { - panic(err) - } - return &Command{ - Server: s, Config: NewConfig(), CmdIO: pilosa.NewCmdIO(stdin, stdout, stderr), @@ -97,7 +91,7 @@ func NewCommand(stdin io.Reader, stdout, stderr io.Writer) *Command { } // Run executes the pilosa server. -func (m *Command) Run(args ...string) (err error) { +func (m *Command) Run(args ...string) (err error) { // TODO args WTF defer close(m.Started) prefix := "~" + string(filepath.Separator) if strings.HasPrefix(m.Config.DataDir, prefix) { @@ -125,76 +119,72 @@ func (m *Command) Run(args ...string) (err error) { return fmt.Errorf("server.Open: %v", err) } - m.Server.Logger.Printf("Listening as %s\n", m.Server.URI) + m.logger.Printf("Listening as %s\n", m.Server.URI) + return nil } +// Wait waits for the server to be closed or interrupted. +func (m *Command) Wait() error { + // First SIGKILL causes server to shut down gracefully. + c := make(chan os.Signal, 2) + signal.Notify(c, os.Interrupt, syscall.SIGTERM) + select { + case sig := <-c: + m.logger.Printf("Received %s; gracefully shutting down...\n", sig.String()) + + // Second signal causes a hard shutdown. + go func() { <-c; os.Exit(1) }() + return errors.Wrap(m.Close(), "closing command") + case <-m.Done: + m.logger.Printf("Server closed externally") + return nil + } +} + // SetupLogger sets up the logger based on the configuration. -func (m *Command) SetupLogger() error { +func (m *Command) SetupLogger() (pilosa.Logger, error) { + if m.logger != nil { + return m.logger, nil + } var err error if m.Config.LogPath == "" { m.logOutput = m.Stderr } else { m.logOutput, err = os.OpenFile(m.Config.LogPath, os.O_RDWR|os.O_CREATE|os.O_APPEND, 0600) if err != nil { - return err + return nil, errors.Wrap(err, "opening file") } } if m.Config.Verbose { - vbl := pilosa.NewVerboseLogger(m.logOutput) - m.logger = vbl.Logger() - m.Server.Logger = vbl + m.logger = pilosa.NewVerboseLogger(m.logOutput) } else { - sl := pilosa.NewStandardLogger(m.logOutput) - m.logger = sl.Logger() - m.Server.Logger = sl + m.logger = pilosa.NewStandardLogger(m.logOutput) } - return nil + return m.logger, nil } // SetupServer uses the cluster configuration to set up this server. func (m *Command) SetupServer() error { - m.Server.Handler.Logger = m.Server.Logger - m.Server.Holder.Logger = m.Server.Logger - m.Server.Holder.Stats.SetLogger(m.Server.Logger) + if m.logger == nil { + _, err := m.SetupLogger() + if err != nil { + return errors.Wrap(err, "setting up logger") + } + + } + + handler := pilosa.NewHandler() + handler.Logger = m.logger + handler.FileSystem = &statik.FileSystem{} + handler.API = pilosa.NewAPI() + handler.API.Logger = m.logger uri, err := pilosa.AddressWithDefaults(m.Config.Bind) - if err != nil { - return err + return errors.Wrap(err, "processing bind address") } - m.Server.URI = *uri - - cluster := pilosa.NewCluster() - cluster.ReplicaN = m.Config.Cluster.ReplicaN - cluster.Holder = m.Server.Holder - cluster.Logger = m.Server.Logger - - m.Server.Cluster = cluster - - // Configure data directory (for Cluster .topology) - m.Server.Cluster.Path = m.Config.DataDir - - m.Server.NewAttrStore = boltdb.NewAttrStore - m.Server.Holder.NewAttrStore = boltdb.NewAttrStore - - // Configure holder. - m.Server.Logger.Printf("Using data from: %s\n", m.Config.DataDir) - m.Server.Holder.Path = m.Config.DataDir - m.Server.MetricInterval = time.Duration(m.Config.Metric.PollInterval) - if m.Config.Metric.Diagnostics { - m.Server.DiagnosticInterval = time.Duration(DefaultDiagnosticsInterval) - } - m.Server.SystemInfo = gopsutil.NewSystemInfo() - m.Server.GCNotifier = gcnotify.NewActiveGCNotifier() - m.Server.Holder.Stats, err = NewStatsClient(m.Config.Metric.Service, m.Config.Metric.Host) - if err != nil { - return err - } - - // Copy configuration flags. - m.Server.MaxWritesPerRequest = m.Config.MaxWritesPerRequest // Setup TLS var TLSConfig *tls.Config @@ -207,27 +197,73 @@ func (m *Command) SetupServer() error { } cert, err := tls.LoadX509KeyPair(m.Config.TLS.CertificatePath, m.Config.TLS.CertificateKeyPath) if err != nil { - return err + return errors.Wrap(err, "load x509 key pair") } - m.Server.TLS = &tls.Config{ + TLSConfig = &tls.Config{ Certificates: []tls.Certificate{cert}, InsecureSkipVerify: m.Config.TLS.SkipVerify, } - - TLSConfig = m.Server.TLS } - c := pilosa.GetHTTPClient(TLSConfig) - m.Server.RemoteClient = c - m.Server.Handler.API.RemoteClient = c - m.Server.Cluster.RemoteClient = c - // Statik file system. - m.Server.Handler.FileSystem = &statik.FileSystem{} + diagnosticsInterval := time.Duration(0) + if m.Config.Metric.Diagnostics { + diagnosticsInterval = time.Duration(DefaultDiagnosticsInterval) + } - // Set configuration options. - m.Server.AntiEntropyInterval = time.Duration(m.Config.AntiEntropy.Interval) - m.Server.Cluster.LongQueryTime = time.Duration(m.Config.Cluster.LongQueryTime) - return nil + statsClient, err := NewStatsClient(m.Config.Metric.Service, m.Config.Metric.Host) + if err != nil { + return errors.Wrap(err, "new stats client") + } + + ln, err := getListener(*uri, TLSConfig) + if err != nil { + return errors.Wrap(err, "getting listener") + } + + c := GetHTTPClient(TLSConfig) + handler.API.RemoteClient = c + + m.Server, err = pilosa.NewServer( + pilosa.OptServerAntiEntropyInterval(time.Duration(m.Config.AntiEntropy.Interval)), + pilosa.OptServerLongQueryTime(time.Duration(m.Config.Cluster.LongQueryTime)), + pilosa.OptServerDataDir(m.Config.DataDir), + pilosa.OptServerReplicaN(m.Config.Cluster.ReplicaN), + pilosa.OptServerMaxWritesPerRequest(m.Config.MaxWritesPerRequest), + pilosa.OptServerMetricInterval(time.Duration(m.Config.Metric.PollInterval)), + pilosa.OptServerDiagnosticsInterval(diagnosticsInterval), + + pilosa.OptServerLogger(m.logger), + pilosa.OptServerAttrStoreFunc(boltdb.NewAttrStore), + pilosa.OptServerHandler(handler), + pilosa.OptServerSystemInfo(gopsutil.NewSystemInfo()), + pilosa.OptServerGCNotifier(gcnotify.NewActiveGCNotifier()), + pilosa.OptServerStatsClient(statsClient), + pilosa.OptServerListener(ln), + pilosa.OptServerURI(uri), + pilosa.OptServerRemoteClient(c), + ) + + return errors.Wrap(err, "new server") +} + +func GetHTTPClient(t *tls.Config) *http.Client { + transport := &http.Transport{ + Proxy: http.ProxyFromEnvironment, + DialContext: (&net.Dialer{ + Timeout: 30 * time.Second, + KeepAlive: 30 * time.Second, + DualStack: true, + }).DialContext, + MaxIdleConns: 1000, + MaxIdleConnsPerHost: 200, + IdleConnTimeout: 90 * time.Second, + TLSHandshakeTimeout: 10 * time.Second, + ExpectContinueTimeout: 1 * time.Second, + } + if t != nil { + transport.TLSClientConfig = t + } + return &http.Client{Transport: transport} } // SetupNetworking sets up internode communication based on the configuration. @@ -266,7 +302,7 @@ func (m *Command) SetupNetworking() error { if m.GossipTransport != nil { transport = m.GossipTransport } else { - transport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger) + transport, err = gossip.NewTransport(gossipHost, gossipPort, m.logger.Logger()) if err != nil { return err } @@ -277,8 +313,8 @@ func (m *Command) SetupNetworking() error { m.Server.Cluster.Coordinator = m.Server.NodeID } - m.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m.Server.Logger) - gossipMemberSet, err := gossip.NewGossipMemberSet(m.Server.NodeID, m.Server.URI.Host(), m.Config.Gossip, m.Server, gossip.WithLogger(m.logger), gossip.WithTransport(transport)) + m.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m.logger) + gossipMemberSet, err := gossip.NewGossipMemberSet(m.Server.NodeID, m.Server.URI.Host(), m.Config.Gossip, m.Server, gossip.WithLogger(m.logger.Logger()), gossip.WithTransport(transport)) if err != nil { return err } @@ -318,3 +354,24 @@ func NewStatsClient(name string, host string) (pilosa.StatsClient, error) { return nil, errors.Errorf("'%v' not a valid stats client, choose from [expvar, statsd, none].") } } + +// OpenListener opens a listener for the Server. +func getListener(uri pilosa.URI, tlsconf *tls.Config) (ln net.Listener, err error) { + // If bind URI has the https scheme, enable TLS + if uri.Scheme() == "https" && tlsconf != nil { + ln, err = tls.Listen("tcp", uri.HostPort(), tlsconf) + if err != nil { + return nil, errors.Wrap(err, "tls.Listener") + } + } else if uri.Scheme() == "http" { + // Open HTTP listener to determine port (if specified as :0). + ln, err = net.Listen("tcp", uri.HostPort()) + if err != nil { + return nil, errors.Wrap(err, "net.Listen") + } + } else { + return nil, errors.Errorf("unsupported scheme: %s", uri.Scheme()) + } + + return ln, nil +} diff --git a/server/server_test.go b/server/server_test.go index 532257c13..75e0e2242 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -45,7 +45,7 @@ func TestMain_Set_Quick(t *testing.T) { defer m.Close() // Create client. - client, err := pilosa.NewInternalHTTPClient(m.Server.URI.HostPort(), pilosa.GetHTTPClient(nil)) + client, err := pilosa.NewInternalHTTPClient(m.Server.URI.HostPort(), server.GetHTTPClient(nil)) if err != nil { t.Fatal(err) } @@ -322,11 +322,11 @@ func TestMain_FrameRestore(t *testing.T) { defer m21.Close() // Import from first cluster. - client20, err := pilosa.NewInternalHTTPClient(m20.Server.URI.HostPort(), pilosa.GetHTTPClient(nil)) + client20, err := pilosa.NewInternalHTTPClient(m20.Server.URI.HostPort(), server.GetHTTPClient(nil)) if err != nil { t.Fatal("new client:", err) } - client21, err := pilosa.NewInternalHTTPClient(m21.Server.URI.HostPort(), pilosa.GetHTTPClient(nil)) + client21, err := pilosa.NewInternalHTTPClient(m21.Server.URI.HostPort(), server.GetHTTPClient(nil)) if err != nil { t.Fatal("new client:", err) } diff --git a/test/executor.go b/test/executor.go index afe244ad7..8cd6391c3 100644 --- a/test/executor.go +++ b/test/executor.go @@ -20,6 +20,7 @@ import ( "github.com/pilosa/pilosa" "github.com/pilosa/pilosa/pql" + "github.com/pilosa/pilosa/server" ) // Executor represents a test wrapper for pilosa.Executor. @@ -30,7 +31,7 @@ type Executor struct { var remoteClient *http.Client func init() { - remoteClient = pilosa.GetHTTPClient(nil) + remoteClient = server.GetHTTPClient(nil) } // NewExecutor returns a new instance of Executor. diff --git a/test/pilosa.go b/test/pilosa.go index b00c9841c..366facdd4 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -49,15 +49,13 @@ func NewMain() *Main { } m := &Main{Command: server.NewCommand(os.Stdin, os.Stdout, os.Stderr)} - m.Server.Network = *Network - m.Server.NewAttrStore = NewAttrStore - m.Server.Holder.NewAttrStore = NewAttrStore m.Config.DataDir = path m.Config.Bind = "http://localhost:0" m.Config.Cluster.Disabled = true m.Command.Stdin = &m.Stdin m.Command.Stdout = &m.Stdout m.Command.Stderr = &m.Stderr + m.SetupServer() if testing.Verbose() { m.Command.Stdout = io.MultiWriter(os.Stdout, m.Command.Stdout) @@ -101,6 +99,7 @@ func runMainWithCluster(size int) ([]*Main, error) { for i := 0; i < size; i++ { m := NewMainWithCluster(i == 0) + m.Config.Cluster.Disabled = false gossipSeeds[i], err = m.RunWithTransport(gossipHost, gossipPort, gossipSeeds[:i]) if err != nil { @@ -136,12 +135,16 @@ func (m *Main) Reopen() error { } // Create new main with the same config. - config := m.Config + config := m.Command.Config m.Command = server.NewCommand(os.Stdin, os.Stdout, os.Stderr) - m.Server.Network = *Network + m.Command.Config = config + err := m.SetupServer() + if err != nil { + return errors.Wrap(err, "setting up server") + } + m.Server.NewAttrStore = boltdb.NewAttrStore m.Server.Holder.NewAttrStore = m.Server.NewAttrStore - m.Config = config // Run new program. if err := m.Run(); err != nil { @@ -174,12 +177,6 @@ func (m *Main) RunWithTransport(host string, bindPort int, joinSeeds []string) ( return seed, err } - // Open server listener. - err = m.Server.OpenListener() - if err != nil { - return seed, err - } - // Open gossip transport to use in SetupServer. transport, err := gossip.NewTransport(host, bindPort, nil) if err != nil { @@ -221,7 +218,7 @@ func (m *Main) URL() string { return "http://" + m.Server.Addr().String() } // Client returns a client to connect to the program. func (m *Main) Client() *pilosa.InternalHTTPClient { - client, err := pilosa.NewInternalHTTPClient(m.Server.URI.HostPort(), pilosa.GetHTTPClient(nil)) + client, err := pilosa.NewInternalHTTPClient(m.Server.URI.HostPort(), server.GetHTTPClient(nil)) if err != nil { panic(err) } diff --git a/test/test.go b/test/test.go index c83b657c0..01f980ad2 100644 --- a/test/test.go +++ b/test/test.go @@ -13,10 +13,3 @@ // limitations under the License. package test - -import "flag" - -// Test flags. -var ( - Network = flag.String("network", "tcp", "network name") -) From 9f1720f01d3e73989f7dde42b9c187bf25486b36 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Mon, 23 Apr 2018 15:12:23 -0500 Subject: [PATCH 2/5] unexport stuff in pilosa.Server refactor gossip.NewGossipMemberset to not take Server --- diagnostics.go | 16 ++--- gossip/gossip.go | 15 ++-- server.go | 182 +++++++++++++++++++++++------------------------ server/server.go | 5 +- 4 files changed, 107 insertions(+), 111 deletions(-) diff --git a/diagnostics.go b/diagnostics.go index 57423d0a5..bfebbb495 100644 --- a/diagnostics.go +++ b/diagnostics.go @@ -168,23 +168,23 @@ func (d *DiagnosticsCollector) logErr(err error) bool { // EnrichWithOSInfo adds OS information to the diagnostics payload. func (d *DiagnosticsCollector) EnrichWithOSInfo() { - uptime, err := d.server.SystemInfo.Uptime() + uptime, err := d.server.systemInfo.Uptime() if !d.logErr(err) { d.Set("HostUptime", uptime) } - platform, err := d.server.SystemInfo.Platform() + platform, err := d.server.systemInfo.Platform() if !d.logErr(err) { d.Set("OSPlatform", platform) } - family, err := d.server.SystemInfo.Family() + family, err := d.server.systemInfo.Family() if !d.logErr(err) { d.Set("OSFamily", family) } - version, err := d.server.SystemInfo.OSVersion() + version, err := d.server.systemInfo.OSVersion() if !d.logErr(err) { d.Set("OSVersion", version) } - kernelVersion, err := d.server.SystemInfo.KernelVersion() + kernelVersion, err := d.server.systemInfo.KernelVersion() if !d.logErr(err) { d.Set("OSKernelVersion", kernelVersion) } @@ -192,15 +192,15 @@ func (d *DiagnosticsCollector) EnrichWithOSInfo() { // EnrichWithMemoryInfo adds memory information to the diagnostics payload. func (d *DiagnosticsCollector) EnrichWithMemoryInfo() { - memFree, err := d.server.SystemInfo.MemFree() + memFree, err := d.server.systemInfo.MemFree() if !d.logErr(err) { d.Set("MemFree", memFree) } - memTotal, err := d.server.SystemInfo.MemTotal() + memTotal, err := d.server.systemInfo.MemTotal() if !d.logErr(err) { d.Set("MemTotal", memTotal) } - memUsed, err := d.server.SystemInfo.MemUsed() + memUsed, err := d.server.systemInfo.MemUsed() if !d.logErr(err) { d.Set("MemUsed", memUsed) } diff --git a/gossip/gossip.go b/gossip/gossip.go index 8ad00a8a7..e4856ec7b 100644 --- a/gossip/gossip.go +++ b/gossip/gossip.go @@ -152,7 +152,7 @@ type gossipConfig struct { type GossipMemberSetOption func(*GossipMemberSet) error // WithTransport is a functional option for providing a transport to NewGossipMemberSet. -func WithTransport(transport *Transport) func(*GossipMemberSet) error { +func WithTransport(transport *Transport) GossipMemberSetOption { return func(g *GossipMemberSet) error { g.transport = transport return nil @@ -160,7 +160,7 @@ func WithTransport(transport *Transport) func(*GossipMemberSet) error { } // WithLogger is a functional option for providing a logger to NewGossipMemberSet. -func WithLogger(logger *log.Logger) func(*GossipMemberSet) error { +func WithLogger(logger *log.Logger) GossipMemberSetOption { return func(g *GossipMemberSet) error { g.logger = logger return nil @@ -168,11 +168,8 @@ func WithLogger(logger *log.Logger) func(*GossipMemberSet) error { } // NewGossipMemberSet returns a new instance of GossipMemberSet based on options. -func NewGossipMemberSet(name string, host string, cfg Config, server *pilosa.Server, options ...GossipMemberSetOption) (*GossipMemberSet, error) { - - g := &GossipMemberSet{ - Logger: server.Logger, - } +func NewGossipMemberSet(name string, host string, cfg Config, ger *GossipEventReceiver, sh pilosa.StatusHandler, options ...GossipMemberSetOption) (*GossipMemberSet, error) { + g := &GossipMemberSet{} // options for _, opt := range options { @@ -227,7 +224,7 @@ func NewGossipMemberSet(name string, host string, cfg Config, server *pilosa.Ser // conf.Delegate = g conf.SecretKey = gossipKey - conf.Events = server.Cluster.EventReceiver.(memberlist.EventDelegate) + conf.Events = ger conf.Logger = g.logger g.config = &gossipConfig{ @@ -235,7 +232,7 @@ func NewGossipMemberSet(name string, host string, cfg Config, server *pilosa.Ser gossipSeeds: cfg.Seeds, } - g.statusHandler = server + g.statusHandler = sh return g, nil } diff --git a/server.go b/server.go index f022d562a..d302a7d97 100644 --- a/server.go +++ b/server.go @@ -16,7 +16,6 @@ package pilosa import ( "context" - "crypto/tls" "fmt" "net" "net/http" @@ -45,12 +44,45 @@ var _ Broadcaster = &Server{} var _ BroadcastHandler = &Server{} var _ StatusHandler = &Server{} +// Server represents a holder wrapped by a running HTTP server. +type Server struct { + // Close management. + wg sync.WaitGroup + closing chan struct{} + + // Internal + Holder *Holder + Cluster *Cluster + diagnostics *DiagnosticsCollector + + // External + handler *Handler + Broadcaster Broadcaster + BroadcastReceiver BroadcastReceiver + Gossiper Gossiper + remoteClient *http.Client + systemInfo SystemInfo + gcNotifier GCNotifier + NewAttrStore func(string) AttrStore + logger Logger + ln net.Listener + + NodeID string + URI URI + antiEntropyInterval time.Duration + metricInterval time.Duration + diagnosticInterval time.Duration + maxWritesPerRequest int + + defaultClient InternalClient +} + // ServerOption is a functional option type for pilosa.Server type ServerOption func(s *Server) error func OptServerLogger(l Logger) ServerOption { return func(s *Server) error { - s.Logger = l + s.logger = l return nil } } @@ -80,7 +112,7 @@ func OptServerAttrStoreFunc(af func(string) AttrStore) ServerOption { func OptServerAntiEntropyInterval(interval time.Duration) ServerOption { return func(s *Server) error { - s.AntiEntropyInterval = interval + s.antiEntropyInterval = interval return nil } } @@ -94,42 +126,42 @@ func OptServerLongQueryTime(dur time.Duration) ServerOption { func OptServerHandler(h *Handler) ServerOption { return func(s *Server) error { - s.Handler = h + s.handler = h return nil } } func OptServerMaxWritesPerRequest(n int) ServerOption { return func(s *Server) error { - s.MaxWritesPerRequest = n + s.maxWritesPerRequest = n return nil } } func OptServerMetricInterval(dur time.Duration) ServerOption { return func(s *Server) error { - s.MetricInterval = dur + s.metricInterval = dur return nil } } func OptServerSystemInfo(si SystemInfo) ServerOption { return func(s *Server) error { - s.SystemInfo = si + s.systemInfo = si return nil } } func OptServerGCNotifier(gcn GCNotifier) ServerOption { return func(s *Server) error { - s.GCNotifier = gcn + s.gcNotifier = gcn return nil } } func OptServerRemoteClient(c *http.Client) ServerOption { return func(s *Server) error { - s.RemoteClient = c + s.remoteClient = c s.Cluster.RemoteClient = c return nil } @@ -144,7 +176,7 @@ func OptServerStatsClient(sc StatsClient) ServerOption { func OptServerDiagnosticsInterval(dur time.Duration) ServerOption { return func(s *Server) error { - s.DiagnosticInterval = dur + s.diagnosticInterval = dur return nil } } @@ -164,61 +196,27 @@ func OptServerURI(uri *URI) ServerOption { } } -// Server represents a holder wrapped by a running HTTP server. -type Server struct { - // Close management. - wg sync.WaitGroup - closing chan struct{} - - // Internal - Holder *Holder - Cluster *Cluster - diagnostics *DiagnosticsCollector - - // External - Handler *Handler - Broadcaster Broadcaster - BroadcastReceiver BroadcastReceiver - Gossiper Gossiper - RemoteClient *http.Client - SystemInfo SystemInfo - GCNotifier GCNotifier - NewAttrStore func(string) AttrStore - Logger Logger - TLS *tls.Config - ln net.Listener - - NodeID string - URI URI - AntiEntropyInterval time.Duration - MetricInterval time.Duration - DiagnosticInterval time.Duration - MaxWritesPerRequest int - - defaultClient InternalClient -} - // NewServer returns a new instance of Server. func NewServer(opts ...ServerOption) (*Server, error) { s := &Server{ closing: make(chan struct{}), Cluster: NewCluster(), Holder: NewHolder(), - Handler: NewHandler(), + handler: NewHandler(), Broadcaster: NopBroadcaster, BroadcastReceiver: NopBroadcastReceiver, diagnostics: NewDiagnosticsCollector(DefaultDiagnosticServer), - SystemInfo: NewNopSystemInfo(), + systemInfo: NewNopSystemInfo(), - GCNotifier: NopGCNotifier, + gcNotifier: NopGCNotifier, NewAttrStore: NewNopAttrStore, - AntiEntropyInterval: time.Minute * 10, - MetricInterval: 0, - DiagnosticInterval: 0, + antiEntropyInterval: time.Minute * 10, + metricInterval: 0, + diagnosticInterval: 0, - Logger: NopLogger, + logger: NopLogger, } for _, opt := range opts { @@ -228,12 +226,12 @@ func NewServer(opts ...ServerOption) (*Server, error) { } } - s.Holder.Logger = s.Logger - s.Holder.Stats.SetLogger(s.Logger) + s.Holder.Logger = s.logger + s.Holder.Stats.SetLogger(s.logger) - s.Cluster.Logger = s.Logger + s.Cluster.Logger = s.logger s.Cluster.Holder = s.Holder - s.Cluster.RemoteClient = s.RemoteClient + s.Cluster.RemoteClient = s.remoteClient // update URI port with actual listener port. TODO this should probably be done outside of here. if s.URI.Port() == 0 { s.URI.SetPort(uint16(s.ln.Addr().(*net.TCPAddr).Port)) @@ -243,7 +241,7 @@ func NewServer(opts ...ServerOption) (*Server, error) { // Open opens and initializes the server. func (s *Server) Open() error { - s.Logger.Printf("open server") + s.logger.Printf("open server") // s.ln can be configured prior to Open() via s.OpenListener(). if s.ln == nil { return errors.New("Must pass a listener option to NewServer") @@ -269,36 +267,36 @@ func (s *Server) Open() error { s.Holder.Peek() // Create default HTTP client - s.createDefaultClient(s.RemoteClient) + s.createDefaultClient(s.remoteClient) // Create executor for executing queries. - e := NewExecutor(s.RemoteClient) + e := NewExecutor(s.remoteClient) e.Holder = s.Holder e.Node = node e.Cluster = s.Cluster - e.MaxWritesPerRequest = s.MaxWritesPerRequest + e.MaxWritesPerRequest = s.maxWritesPerRequest // Cluster settings. s.Cluster.Broadcaster = s.Broadcaster - s.Cluster.MaxWritesPerRequest = s.MaxWritesPerRequest + s.Cluster.MaxWritesPerRequest = s.maxWritesPerRequest // Initialize HTTP handler. - s.Handler.API.Holder = s.Holder - s.Handler.API.Broadcaster = s.Broadcaster - s.Handler.API.BroadcastHandler = s - s.Handler.API.StatusHandler = s - s.Handler.API.URI = s.URI - s.Handler.API.Cluster = s.Cluster - s.Handler.API.Executor = e + s.handler.API.Holder = s.Holder + s.handler.API.Broadcaster = s.Broadcaster + s.handler.API.BroadcastHandler = s + s.handler.API.StatusHandler = s + s.handler.API.URI = s.URI + s.handler.API.Cluster = s.Cluster + s.handler.API.Executor = e // Initialize Holder. s.Holder.Broadcaster = s.Broadcaster // Serve HTTP. go func() { - err := http.Serve(s.ln, s.Handler) + err := http.Serve(s.ln, s.handler) if err != nil { - s.Logger.Printf("HTTP handler terminated with error: %s\n", err) + s.logger.Printf("HTTP handler terminated with error: %s\n", err) } }() @@ -363,7 +361,7 @@ func (s *Server) LoadNodeID() string { } nodeID, err := s.Holder.loadNodeID() if err != nil { - s.Logger.Printf("loading NodeID: %v", err) + s.logger.Printf("loading NodeID: %v", err) return s.NodeID } return nodeID @@ -378,10 +376,10 @@ func (s *Server) Addr() net.Addr { } func (s *Server) monitorAntiEntropy() { - ticker := time.NewTicker(s.AntiEntropyInterval) + ticker := time.NewTicker(s.antiEntropyInterval) defer ticker.Stop() - s.Logger.Printf("holder sync monitor initializing (%s interval)", s.AntiEntropyInterval) + s.logger.Printf("holder sync monitor initializing (%s interval)", s.antiEntropyInterval) for { // Wait for tick or a close. @@ -392,7 +390,7 @@ func (s *Server) monitorAntiEntropy() { s.Holder.Stats.Count("AntiEntropy", 1, 1.0) } t := time.Now() - s.Logger.Printf("holder sync beginning") + s.logger.Printf("holder sync beginning") // Initialize syncer with local holder and remote client. var syncer HolderSyncer @@ -400,17 +398,17 @@ func (s *Server) monitorAntiEntropy() { syncer.Node = s.Cluster.Node syncer.Cluster = s.Cluster syncer.Closing = s.closing - syncer.RemoteClient = s.RemoteClient + syncer.RemoteClient = s.remoteClient syncer.Stats = s.Holder.Stats.WithTags("HolderSyncer") // Sync holders. if err := syncer.SyncHolder(); err != nil { - s.Logger.Printf("holder sync error: err=%s", err) + s.logger.Printf("holder sync error: err=%s", err) continue } // Record successful sync in log. - s.Logger.Printf("holder sync complete") + s.logger.Printf("holder sync complete") dif := time.Since(t) s.Holder.Stats.Histogram("AntiEntropyDuration", float64(dif), 1.0) } @@ -532,7 +530,7 @@ func (s *Server) ReceiveMessage(pb proto.Message) error { func (s *Server) SendSync(pb proto.Message) error { var eg errgroup.Group for _, node := range s.Cluster.Nodes { - s.Logger.Printf("SendSync to: %s", node.URI) + s.logger.Printf("SendSync to: %s", node.URI) // Don't forward the message to ourselves. if s.URI == node.URI { continue @@ -554,7 +552,7 @@ func (s *Server) SendAsync(pb proto.Message) error { // SendTo represents an implementation of Broadcaster. func (s *Server) SendTo(to *Node, pb proto.Message) error { - s.Logger.Printf("SendTo: %s", to.URI) + s.logger.Printf("SendTo: %s", to.URI) ctx := context.WithValue(context.Background(), "uri", &to.URI) return s.defaultClient.SendMessage(ctx, pb) } @@ -604,7 +602,7 @@ func (s *Server) HandleRemoteStatus(pb proto.Message) error { err := s.mergeRemoteStatus(pb.(*internal.NodeStatus)) if err != nil { - s.Logger.Printf("merge remote status: %s", err) + s.logger.Printf("merge remote status: %s", err) } }() @@ -629,7 +627,7 @@ func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error { // if we don't know about an index locally, log an error because // indexes should be created and synced prior to slice creation if localIndex == nil { - s.Logger.Printf("Local Index not found: %s", index) + s.logger.Printf("Local Index not found: %s", index) continue } if newMax > oldmaxslices[index] { @@ -645,7 +643,7 @@ func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error { // if we don't know about an index locally, log an error because // indexes should be created and synced prior to slice creation if localIndex == nil { - s.Logger.Printf("Local Index not found: %s", index) + s.logger.Printf("Local Index not found: %s", index) continue } if newMaxInverse > oldMaxInverseSlices[index] { @@ -660,14 +658,14 @@ func (s *Server) mergeRemoteStatus(ns *internal.NodeStatus) error { // monitorDiagnostics periodically polls the Pilosa Indexes for cluster info. func (s *Server) monitorDiagnostics() { // Do not send more than once a minute - if s.DiagnosticInterval < time.Minute { - s.Logger.Printf("diagnostics disabled") + if s.diagnosticInterval < time.Minute { + s.logger.Printf("diagnostics disabled") return } else { - s.Logger.Printf("Pilosa is currently configured to send small diagnostics reports to our team every %v. More information here: https://www.pilosa.com/docs/latest/administration/#diagnostics", s.DiagnosticInterval) + s.logger.Printf("Pilosa is currently configured to send small diagnostics reports to our team every %v. More information here: https://www.pilosa.com/docs/latest/administration/#diagnostics", s.diagnosticInterval) } - s.diagnostics.Logger = s.Logger + s.diagnostics.Logger = s.logger s.diagnostics.SetVersion(Version) s.diagnostics.Set("Host", s.URI.host) s.diagnostics.Set("Cluster", strings.Join(s.Cluster.NodeIDs(), ",")) @@ -689,11 +687,11 @@ func (s *Server) monitorDiagnostics() { s.diagnostics.CheckVersion() err = s.diagnostics.Flush() if err != nil { - s.Logger.Printf("Diagnostics error: %s", err) + s.logger.Printf("Diagnostics error: %s", err) } } - ticker := time.NewTicker(s.DiagnosticInterval) + ticker := time.NewTicker(s.diagnosticInterval) defer ticker.Stop() flush() for { @@ -710,24 +708,24 @@ func (s *Server) monitorDiagnostics() { // monitorRuntime periodically polls the Go runtime metrics. func (s *Server) monitorRuntime() { // Disable metrics when poll interval is zero. - if s.MetricInterval <= 0 { + if s.metricInterval <= 0 { return } var m runtime.MemStats - ticker := time.NewTicker(s.MetricInterval) + ticker := time.NewTicker(s.metricInterval) defer ticker.Stop() - defer s.GCNotifier.Close() + defer s.gcNotifier.Close() - s.Logger.Printf("runtime stats initializing (%s interval)", s.MetricInterval) + s.logger.Printf("runtime stats initializing (%s interval)", s.metricInterval) for { // Wait for tick or a close. select { case <-s.closing: return - case <-s.GCNotifier.AfterGC(): + case <-s.gcNotifier.AfterGC(): // GC just ran. s.Holder.Stats.Count("garbage_collection", 1, 1.0) case <-ticker.C: diff --git a/server/server.go b/server/server.go index 829676fd5..3e5bb1899 100644 --- a/server/server.go +++ b/server/server.go @@ -313,8 +313,9 @@ func (m *Command) SetupNetworking() error { m.Server.Cluster.Coordinator = m.Server.NodeID } - m.Server.Cluster.EventReceiver = gossip.NewGossipEventReceiver(m.logger) - gossipMemberSet, err := gossip.NewGossipMemberSet(m.Server.NodeID, m.Server.URI.Host(), m.Config.Gossip, m.Server, gossip.WithLogger(m.logger.Logger()), gossip.WithTransport(transport)) + gossipEventReceiver := gossip.NewGossipEventReceiver(m.logger) + m.Server.Cluster.EventReceiver = gossipEventReceiver + gossipMemberSet, err := gossip.NewGossipMemberSet(m.Server.NodeID, m.Server.URI.Host(), m.Config.Gossip, gossipEventReceiver, m.Server, gossip.WithLogger(m.logger.Logger()), gossip.WithTransport(transport)) if err != nil { return err } From 1e05a6d62725b0e261f09e32ddd4322b6a6ace34 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Mon, 23 Apr 2018 15:15:02 -0500 Subject: [PATCH 3/5] fix getListener comment --- server/server.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/server/server.go b/server/server.go index 3e5bb1899..1c3ba577c 100644 --- a/server/server.go +++ b/server/server.go @@ -356,7 +356,7 @@ func NewStatsClient(name string, host string) (pilosa.StatsClient, error) { } } -// OpenListener opens a listener for the Server. +// getListener gets a net.Listener based on the config. func getListener(uri pilosa.URI, tlsconf *tls.Config) (ln net.Listener, err error) { // If bind URI has the https scheme, enable TLS if uri.Scheme() == "https" && tlsconf != nil { From 9deefbaefc00366008b2ef25aa7cb880b3497da3 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Mon, 23 Apr 2018 16:25:03 -0500 Subject: [PATCH 4/5] unexport setupLogger and simplify --- server/server.go | 20 +++++++------------- 1 file changed, 7 insertions(+), 13 deletions(-) diff --git a/server/server.go b/server/server.go index 1c3ba577c..6ff80fe7f 100644 --- a/server/server.go +++ b/server/server.go @@ -142,18 +142,15 @@ func (m *Command) Wait() error { } } -// SetupLogger sets up the logger based on the configuration. -func (m *Command) SetupLogger() (pilosa.Logger, error) { - if m.logger != nil { - return m.logger, nil - } +// setupLogger sets up the logger based on the configuration. +func (m *Command) setupLogger() error { var err error if m.Config.LogPath == "" { m.logOutput = m.Stderr } else { m.logOutput, err = os.OpenFile(m.Config.LogPath, os.O_RDWR|os.O_CREATE|os.O_APPEND, 0600) if err != nil { - return nil, errors.Wrap(err, "opening file") + return errors.Wrap(err, "opening file") } } @@ -162,17 +159,14 @@ func (m *Command) SetupLogger() (pilosa.Logger, error) { } else { m.logger = pilosa.NewStandardLogger(m.logOutput) } - return m.logger, nil + return nil } // SetupServer uses the cluster configuration to set up this server. func (m *Command) SetupServer() error { - if m.logger == nil { - _, err := m.SetupLogger() - if err != nil { - return errors.Wrap(err, "setting up logger") - } - + err := m.setupLogger() + if err != nil { + return errors.Wrap(err, "setting up logger") } handler := pilosa.NewHandler() From f99479932db488f42c3f61efc37eff5d7c0608fd Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Mon, 23 Apr 2018 16:28:24 -0500 Subject: [PATCH 5/5] rename server Run to Start to better reflect functionality --- cmd/server.go | 2 +- server/server.go | 4 ++-- test/pilosa.go | 4 ++-- 3 files changed, 5 insertions(+), 5 deletions(-) diff --git a/cmd/server.go b/cmd/server.go index 93ecb543b..d906cf189 100644 --- a/cmd/server.go +++ b/cmd/server.go @@ -39,7 +39,7 @@ It will load existing data from the configured directory and start listening for client connections on the configured port.`, RunE: func(cmd *cobra.Command, args []string) error { - if err := Server.Run(); err != nil { + if err := Server.Start(); err != nil { return errors.Wrap(err, "running server") } return errors.Wrap(Server.Wait(), "waiting on Server") diff --git a/server/server.go b/server/server.go index 6ff80fe7f..7f5782211 100644 --- a/server/server.go +++ b/server/server.go @@ -90,8 +90,8 @@ func NewCommand(stdin io.Reader, stdout, stderr io.Writer) *Command { } } -// Run executes the pilosa server. -func (m *Command) Run(args ...string) (err error) { // TODO args WTF +// Start starts the pilosa server - it returns once the server is running. +func (m *Command) Start() (err error) { defer close(m.Started) prefix := "~" + string(filepath.Separator) if strings.HasPrefix(m.Config.DataDir, prefix) { diff --git a/test/pilosa.go b/test/pilosa.go index 366facdd4..a390cd462 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -116,7 +116,7 @@ func runMainWithCluster(size int) ([]*Main, error) { func MustRunMain() *Main { m := NewMain() m.Config.Metric.Diagnostics = false // Disable diagnostics. - if err := m.Run(); err != nil { + if err := m.Start(); err != nil { panic(err) } return m @@ -147,7 +147,7 @@ func (m *Main) Reopen() error { m.Server.Holder.NewAttrStore = m.Server.NewAttrStore // Run new program. - if err := m.Run(); err != nil { + if err := m.Start(); err != nil { return err } return nil