From 9f1720f01d3e73989f7dde42b9c187bf25486b36 Mon Sep 17 00:00:00 2001 From: Matthew Jaffee Date: Mon, 23 Apr 2018 15:12:23 -0500 Subject: [PATCH] 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 }