unexport stuff in pilosa.Server

refactor gossip.NewGossipMemberset to not take Server
This commit is contained in:
Matthew Jaffee 2018-04-23 15:12:23 -05:00
parent 28acc29a10
commit 9f1720f01d
No known key found for this signature in database
GPG key ID: 51C676AF9FFCDB87
4 changed files with 107 additions and 111 deletions

View file

@ -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)
}

View file

@ -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
}

182
server.go
View file

@ -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:

View file

@ -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
}