diff --git a/client/rpc_client.go b/client/rpc_client.go index 6173cc00e..6a5becf1c 100644 --- a/client/rpc_client.go +++ b/client/rpc_client.go @@ -202,8 +202,8 @@ func (c *RPCClient) ForceLeave(node string) error { return c.genericRPC(&header, &req, nil) } -//ForceLeavePrune uses ForceLeave but is used to reap the -//node entirely +// ForceLeavePrune uses ForceLeave but is used to reap the +// node entirely func (c *RPCClient) ForceLeavePrune(node string) error { header := requestHeader{ Command: forceLeaveCommand, diff --git a/cmd/serf/command/agent/agent.go b/cmd/serf/command/agent/agent.go index e9c319b0f..a12468841 100644 --- a/cmd/serf/command/agent/agent.go +++ b/cmd/serf/command/agent/agent.go @@ -9,11 +9,11 @@ import ( "fmt" "io" "io/ioutil" - "log" "os" "strings" "sync" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/memberlist" "github.com/hashicorp/serf/serf" ) @@ -37,7 +37,7 @@ type Agent struct { eventHandlersLock sync.Mutex // logger instance wraps the logOutput - logger *log.Logger + logger hclog.Logger // This is the underlying Serf we are wrapping serf *serf.Serf @@ -70,8 +70,10 @@ func Create(agentConf *Config, conf *serf.Config, logOutput io.Writer) (*Agent, agentConf: agentConf, eventCh: eventCh, eventHandlers: make(map[EventHandler]struct{}), - logger: log.New(logOutput, "", log.LstdFlags), - shutdownCh: make(chan struct{}), + logger: hclog.New(&hclog.LoggerOptions{ + Output: logOutput, + }), + shutdownCh: make(chan struct{}), } // Restore agent tags from a tags file @@ -95,7 +97,7 @@ func Create(agentConf *Config, conf *serf.Config, logOutput io.Writer) (*Agent, // create so that there isn't a race condition between creating the // agent and registering handlers func (a *Agent) Start() error { - a.logger.Printf("[INFO] agent: Serf agent starting") + a.logger.Info("agent: Serf agent starting") // Create serf first serf, err := serf.Create(a.conf) @@ -115,7 +117,7 @@ func (a *Agent) Leave() error { return nil } - a.logger.Println("[INFO] agent: requesting graceful leave from Serf") + a.logger.Info("agent: requesting graceful leave from Serf") return a.serf.Leave() } @@ -133,13 +135,13 @@ func (a *Agent) Shutdown() error { goto EXIT } - a.logger.Println("[INFO] agent: requesting serf shutdown") + a.logger.Info("agent: requesting serf shutdown") if err := a.serf.Shutdown(); err != nil { return err } EXIT: - a.logger.Println("[INFO] agent: shutdown complete") + a.logger.Info("agent: shutdown complete") a.shutdown = true close(a.shutdownCh) return nil @@ -163,24 +165,24 @@ func (a *Agent) SerfConfig() *serf.Config { // Join asks the Serf instance to join. See the Serf.Join function. func (a *Agent) Join(addrs []string, replay bool) (n int, err error) { - a.logger.Printf("[INFO] agent: joining: %v replay: %v", addrs, replay) + a.logger.Info(fmt.Sprintf("agent: joining: %v replay: %v", addrs, replay)) ignoreOld := !replay n, err = a.serf.Join(addrs, ignoreOld) if n > 0 { - a.logger.Printf("[INFO] agent: joined: %d nodes", n) + a.logger.Info(fmt.Sprintf("agent: joined: %d nodes", n)) } if err != nil { - a.logger.Printf("[WARN] agent: error joining: %v", err) + a.logger.Warn(fmt.Sprintf("agent: error joining: %v", err)) } return } // ForceLeave is used to eject a failed node from the cluster func (a *Agent) ForceLeave(node string) error { - a.logger.Printf("[INFO] agent: Force leaving node: %s", node) + a.logger.Info(fmt.Sprintf("agent: Force leaving node: %s", node)) err := a.serf.RemoveFailedNode(node) if err != nil { - a.logger.Printf("[WARN] agent: failed to remove node: %v", err) + a.logger.Warn(fmt.Sprintf("agent: failed to remove node: %v", err)) } return err } @@ -188,21 +190,21 @@ func (a *Agent) ForceLeave(node string) error { // ForceLeavePrune completely removes a failed node from the // member list entirely func (a *Agent) ForceLeavePrune(node string) error { - a.logger.Printf("[INFO] agent: Force leaving node (prune): %s", node) + a.logger.Info(fmt.Sprintf("agent: Force leaving node (prune): %s", node)) err := a.serf.RemoveFailedNodePrune(node) if err != nil { - a.logger.Printf("[WARN] agent: failed to remove node (prune): %v", err) + a.logger.Warn(fmt.Sprintf("agent: failed to remove node (prune): %v", err)) } return err } // UserEvent sends a UserEvent on Serf, see Serf.UserEvent. func (a *Agent) UserEvent(name string, payload []byte, coalesce bool) error { - a.logger.Printf("[DEBUG] agent: Requesting user event send: %s. Coalesced: %#v. Payload: %#v", - name, coalesce, string(payload)) + a.logger.Debug(fmt.Sprintf("agent: Requesting user event send: %s. Coalesced: %#v. Payload: %#v", + name, coalesce, string(payload))) err := a.serf.UserEvent(name, payload, coalesce) if err != nil { - a.logger.Printf("[WARN] agent: failed to send user event: %v", err) + a.logger.Warn("agent: failed to send user event: %v", err) } return err } @@ -216,11 +218,11 @@ func (a *Agent) Query(name string, payload []byte, params *serf.QueryParam) (*se return nil, fmt.Errorf("Queries cannot contain the '%s' prefix", serf.InternalQueryPrefix) } } - a.logger.Printf("[DEBUG] agent: Requesting query send: %s. Payload: %#v", - name, string(payload)) + a.logger.Debug(fmt.Sprintf("agent: Requesting query send: %s. Payload: %#v", + name, string(payload))) resp, err := a.serf.Query(name, payload, params) if err != nil { - a.logger.Printf("[WARN] agent: failed to start user query: %v", err) + a.logger.Warn(fmt.Sprintf("agent: failed to start user query: %v", err)) } return resp, err } @@ -255,7 +257,7 @@ func (a *Agent) eventLoop() { for { select { case e := <-a.eventCh: - a.logger.Printf("[INFO] agent: Received event: %s", e.String()) + a.logger.Info(fmt.Sprintf("agent: Received event: %s", e.String())) a.eventHandlersLock.Lock() handlers := a.eventHandlerList a.eventHandlersLock.Unlock() @@ -264,7 +266,7 @@ func (a *Agent) eventLoop() { } case <-serfShutdownCh: - a.logger.Printf("[WARN] agent: Serf shutdown detected, quitting") + a.logger.Warn("agent: Serf shutdown detected, quitting") a.Shutdown() return @@ -276,28 +278,28 @@ func (a *Agent) eventLoop() { // InstallKey initiates a query to install a new key on all members func (a *Agent) InstallKey(key string) (*serf.KeyResponse, error) { - a.logger.Print("[INFO] agent: Initiating key installation") + a.logger.Info("agent: Initiating key installation") manager := a.serf.KeyManager() return manager.InstallKey(key) } // UseKey sends a query instructing all members to switch primary keys func (a *Agent) UseKey(key string) (*serf.KeyResponse, error) { - a.logger.Print("[INFO] agent: Initiating primary key change") + a.logger.Info("agent: Initiating primary key change") manager := a.serf.KeyManager() return manager.UseKey(key) } // RemoveKey sends a query to all members to remove a key from the keyring func (a *Agent) RemoveKey(key string) (*serf.KeyResponse, error) { - a.logger.Print("[INFO] agent: Initiating key removal") + a.logger.Info("agent: Initiating key removal") manager := a.serf.KeyManager() return manager.RemoveKey(key) } // ListKeys sends a query to all members to return a list of their keys func (a *Agent) ListKeys() (*serf.KeyResponse, error) { - a.logger.Print("[INFO] agent: Initiating key listing") + a.logger.Info("agent: Initiating key listing") manager := a.serf.KeyManager() return manager.ListKeys() } @@ -308,7 +310,7 @@ func (a *Agent) SetTags(tags map[string]string) error { // Update the tags file if we have one if a.agentConf.TagsFile != "" { if err := a.writeTagsFile(tags); err != nil { - a.logger.Printf("[ERR] agent: %s", err) + a.logger.Error("agent: %s", err) return err } } @@ -333,7 +335,7 @@ func (a *Agent) loadTagsFile(tagsFile string) error { if err := json.Unmarshal(tagData, &a.conf.Tags); err != nil { return fmt.Errorf("Failed to decode tags file: %s", err) } - a.logger.Printf("[INFO] agent: Restored %d tag(s) from %s", + a.logger.Info("agent: Restored %d tag(s) from %s", len(a.conf.Tags), tagsFile) } @@ -425,8 +427,8 @@ func (a *Agent) loadKeyringFile(keyringFile string) error { return fmt.Errorf("Failed to restore keyring: %s", err) } a.conf.MemberlistConfig.Keyring = keyring - a.logger.Printf("[INFO] agent: Restored keyring with %d keys from %s", - len(keys), keyringFile) + a.logger.Info(fmt.Sprintf("agent: Restored keyring with %d keys from %s", + len(keys), keyringFile)) // Success! return nil diff --git a/cmd/serf/command/agent/command.go b/cmd/serf/command/agent/command.go index 3feff1e2c..5f0334a87 100644 --- a/cmd/serf/command/agent/command.go +++ b/cmd/serf/command/agent/command.go @@ -7,7 +7,6 @@ import ( "flag" "fmt" "io" - "log" "net" "os" "os/signal" @@ -17,8 +16,8 @@ import ( "time" "github.com/armon/go-metrics" + "github.com/hashicorp/go-hclog" gsyslog "github.com/hashicorp/go-syslog" - "github.com/hashicorp/logutils" "github.com/hashicorp/memberlist" "github.com/hashicorp/serf/serf" "github.com/mitchellh/cli" @@ -44,8 +43,7 @@ type Command struct { ShutdownCh <-chan struct{} args []string scriptHandler *ScriptEventHandler - logFilter *logutils.LevelFilter - logger *log.Logger + logger hclog.Logger } var _ cli.Command = &Command{} @@ -370,16 +368,6 @@ func (c *Command) setupLoggers(config *Config) (*GatedWriter, *logWriter, io.Wri Writer: &cli.UiWriter{Ui: c.Ui}, } - c.logFilter = LevelFilter() - c.logFilter.MinLevel = logutils.LogLevel(strings.ToUpper(config.LogLevel)) - c.logFilter.Writer = logGate - if !ValidateLevelFilter(c.logFilter.MinLevel, c.logFilter) { - c.Ui.Error(fmt.Sprintf( - "Invalid log level: %s. Valid log levels are: %v", - c.logFilter.MinLevel, c.logFilter.Levels)) - return nil, nil, nil - } - // Check if syslog is enabled var syslog io.Writer if config.EnableSyslog { @@ -388,20 +376,23 @@ func (c *Command) setupLoggers(config *Config) (*GatedWriter, *logWriter, io.Wri c.Ui.Error(fmt.Sprintf("Syslog setup failed: %v", err)) return nil, nil, nil } - syslog = &SyslogWrapper{l, c.logFilter} + syslog = &SyslogWrapper{l} } // Create a log writer, and wrap a logOutput around it logWriter := NewLogWriter(512) var logOutput io.Writer if syslog != nil { - logOutput = io.MultiWriter(c.logFilter, logWriter, syslog) + logOutput = io.MultiWriter(logWriter, syslog) } else { - logOutput = io.MultiWriter(c.logFilter, logWriter) + logOutput = logWriter } // Create a logger - c.logger = log.New(logOutput, "", log.LstdFlags) + c.logger = hclog.New(&hclog.LoggerOptions{ + Output: logOutput, + Level: hclog.LevelFromString(config.LogLevel), + }) return logGate, logWriter, logOutput } @@ -412,7 +403,9 @@ func (c *Command) startAgent(config *Config, agent *Agent, c.scriptHandler = &ScriptEventHandler{ SelfFunc: func() serf.Member { return agent.Serf().LocalMember() }, Scripts: config.EventScripts(), - Logger: log.New(logOutput, "", log.LstdFlags), + Logger: hclog.New(&hclog.LoggerOptions{ + Output: logOutput, + }), } agent.RegisterEventHandler(c.scriptHandler) @@ -504,23 +497,23 @@ func (c *Command) retryJoin(config *Config, agent *Agent, errCh chan struct{}) { attempt := 0 for { // Try to perform the join - c.logger.Printf("[INFO] agent: Joining cluster...(replay: %v)", config.ReplayOnJoin) + c.logger.Info(fmt.Sprintf("agent: Joining cluster...(replay: %v)", config.ReplayOnJoin)) n, err := agent.Join(config.RetryJoin, config.ReplayOnJoin) if err == nil { - c.logger.Printf("[INFO] agent: Join completed. Synced with %d initial agents", n) + c.logger.Info(fmt.Sprintf("agent: Join completed. Synced with %d initial agents", n)) return } // Check if the maximum attempts has been exceeded attempt++ if config.RetryMaxAttempts > 0 && attempt > config.RetryMaxAttempts { - c.logger.Printf("[ERR] agent: maximum retry join attempts made, exiting") + c.logger.Error("agent: maximum retry join attempts made, exiting") close(errCh) return } // Log the failure and sleep - c.logger.Printf("[WARN] agent: Join failed: %v, retrying in %v", err, config.RetryInterval) + c.logger.Warn(fmt.Sprintf("agent: Join failed: %v, retrying in %v", err, config.RetryInterval)) time.Sleep(config.RetryInterval) } } @@ -686,22 +679,11 @@ func (c *Command) handleReload(config *Config, agent *Agent) *Config { c.Ui.Output("Reloading configuration...") newConf := c.readConfig() if newConf == nil { - c.Ui.Error(fmt.Sprintf("Failed to reload configs")) + c.Ui.Error("Failed to reload configs") return config } - // Change the log level - minLevel := logutils.LogLevel(strings.ToUpper(newConf.LogLevel)) - if ValidateLevelFilter(minLevel, c.logFilter) { - c.logFilter.SetMinLevel(minLevel) - } else { - c.Ui.Error(fmt.Sprintf( - "Invalid log level: %s. Valid log levels are: %v", - minLevel, c.logFilter.Levels)) - - // Keep the current log level - newConf.LogLevel = config.LogLevel - } + c.logger.SetLevel(hclog.LevelFromString(newConf.LogLevel)) // Change the event handlers c.scriptHandler.UpdateScripts(newConf.EventScripts()) diff --git a/cmd/serf/command/agent/event_handler.go b/cmd/serf/command/agent/event_handler.go index bbbff31de..6e1343b70 100644 --- a/cmd/serf/command/agent/event_handler.go +++ b/cmd/serf/command/agent/event_handler.go @@ -5,11 +5,11 @@ package agent import ( "fmt" - "log" "os" "strings" "sync" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/serf/serf" ) @@ -22,7 +22,7 @@ type EventHandler interface { type ScriptEventHandler struct { SelfFunc func() serf.Member Scripts []EventScript - Logger *log.Logger + Logger hclog.Logger scriptLock sync.Mutex newScripts []EventScript @@ -38,7 +38,7 @@ func (h *ScriptEventHandler) HandleEvent(e serf.Event) { h.scriptLock.Unlock() if h.Logger == nil { - h.Logger = log.New(os.Stderr, "", log.LstdFlags) + h.Logger = hclog.New(&hclog.LoggerOptions{Output: os.Stderr}) } self := h.SelfFunc() @@ -49,8 +49,8 @@ func (h *ScriptEventHandler) HandleEvent(e serf.Event) { err := invokeEventScript(h.Logger, script.Script, self, e) if err != nil { - h.Logger.Printf("[ERR] agent: Error invoking script '%s': %s", - script.Script, err) + h.Logger.Error(fmt.Sprintf("agent: Error invoking script '%s': %s", + script.Script, err)) } } } diff --git a/cmd/serf/command/agent/invoke.go b/cmd/serf/command/agent/invoke.go index 8a4dbea0a..a4f25c131 100644 --- a/cmd/serf/command/agent/invoke.go +++ b/cmd/serf/command/agent/invoke.go @@ -6,7 +6,6 @@ package agent import ( "fmt" "io" - "log" "os" "os/exec" "regexp" @@ -16,6 +15,7 @@ import ( "github.com/armon/circbuf" "github.com/armon/go-metrics" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/serf/serf" ) @@ -42,7 +42,7 @@ var sanitizeTagRegexp = regexp.MustCompile(`[^A-Z0-9_]`) // // In all events, data is passed in via stdin to facilitate piping. See // the various stdin functions below for more information. -func invokeEventScript(logger *log.Logger, script string, self serf.Member, event serf.Event) error { +func invokeEventScript(logger hclog.Logger, script string, self serf.Member, event serf.Event) error { defer metrics.MeasureSinceWithLabels([]string{"agent", "invoke", script}, time.Now(), nil) output, _ := circbuf.NewBuffer(maxBufSize) @@ -98,8 +98,8 @@ func invokeEventScript(logger *log.Logger, script string, self serf.Member, even // Start a timer to warn about slow handlers slowTimer := time.AfterFunc(warnSlow, func() { - logger.Printf("[WARN] agent: Script '%s' slow, execution exceeding %v", - script, warnSlow) + logger.Warn(fmt.Sprintf("agent: Script '%s' slow, execution exceeding %v", + script, warnSlow)) }) if err := cmd.Start(); err != nil { @@ -108,14 +108,14 @@ func invokeEventScript(logger *log.Logger, script string, self serf.Member, even // Warn if buffer is overritten if output.TotalWritten() > output.Size() { - logger.Printf("[WARN] agent: Script '%s' generated %d bytes of output, truncated to %d", - script, output.TotalWritten(), output.Size()) + logger.Warn(fmt.Sprintf("agent: Script '%s' generated %d bytes of output, truncated to %d", + script, output.TotalWritten(), output.Size())) } err = cmd.Wait() slowTimer.Stop() - logger.Printf("[DEBUG] agent: Event '%s' script output: %s", - event.EventType().String(), output.String()) + logger.Debug(fmt.Sprintf("agent: Event '%s' script output: %s", + event.EventType().String(), output.String())) if err != nil { return err } @@ -123,8 +123,8 @@ func invokeEventScript(logger *log.Logger, script string, self serf.Member, even // If this is a query and we have output, respond if query, ok := event.(*serf.Query); ok && output.TotalWritten() > 0 { if err := query.Respond(output.Bytes()); err != nil { - logger.Printf("[WARN] agent: Failed to respond to query '%s': %s", - event.String(), err) + logger.Warn(fmt.Sprintf("agent: Failed to respond to query '%s': %s", + event.String(), err)) } } @@ -145,7 +145,7 @@ func eventClean(v string) string { // "NAME ADDRESS ROLE TAGS" where the whitespace is actually tabs. // The name and role are cleaned so that newlines and tabs are replaced // with "\n" and "\t" respectively. -func memberEventStdin(logger *log.Logger, stdin io.WriteCloser, e *serf.MemberEvent) { +func memberEventStdin(logger hclog.Logger, stdin io.WriteCloser, e *serf.MemberEvent) { defer stdin.Close() for _, member := range e.Members { // Format the tags as tag1=v1,tag2=v2,... @@ -171,7 +171,7 @@ func memberEventStdin(logger *log.Logger, stdin io.WriteCloser, e *serf.MemberEv // Sends data on stdin for an event. The stdin simply contains the // payload (if any). // Most shells read implementations need a newline, force it to be there -func streamPayload(logger *log.Logger, stdin io.WriteCloser, buf []byte) { +func streamPayload(logger hclog.Logger, stdin io.WriteCloser, buf []byte) { defer stdin.Close() // Append a newline to payload if missing @@ -181,7 +181,7 @@ func streamPayload(logger *log.Logger, stdin io.WriteCloser, buf []byte) { } if _, err := stdin.Write(payload); err != nil { - logger.Printf("[ERR] Error writing payload: %s", err) + logger.Error(fmt.Sprintf("Error writing payload: %s", err)) return } } diff --git a/cmd/serf/command/agent/ipc.go b/cmd/serf/command/agent/ipc.go index c4a532d59..04e05f515 100644 --- a/cmd/serf/command/agent/ipc.go +++ b/cmd/serf/command/agent/ipc.go @@ -28,7 +28,6 @@ import ( "bufio" "fmt" "io" - "log" "net" "os" "regexp" @@ -38,6 +37,7 @@ import ( "time" "github.com/armon/go-metrics" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/go-msgpack/codec" "github.com/hashicorp/logutils" "github.com/hashicorp/serf/coordinate" @@ -245,7 +245,7 @@ type AgentIPC struct { authKey string clients map[string]*IPCClient listener net.Listener - logger *log.Logger + logger hclog.Logger logWriter *logWriter stop uint32 stopCh chan struct{} @@ -335,11 +335,13 @@ func NewAgentIPC(agent *Agent, authKey string, listener net.Listener, logOutput = os.Stderr } ipc := &AgentIPC{ - agent: agent, - authKey: authKey, - clients: make(map[string]*IPCClient), - listener: listener, - logger: log.New(logOutput, "", log.LstdFlags), + agent: agent, + authKey: authKey, + clients: make(map[string]*IPCClient), + listener: listener, + logger: hclog.New(&hclog.LoggerOptions{ + Output: logWriter, + }), logWriter: logWriter, stopCh: make(chan struct{}), } @@ -378,10 +380,10 @@ func (i *AgentIPC) listen() { if i.isStopped() { return } - i.logger.Printf("[ERR] agent.ipc: Failed to accept client: %v", err) + i.logger.Error(fmt.Sprintf("agent.ipc: Failed to accept client: %v", err)) continue } - i.logger.Printf("[INFO] agent.ipc: Accepted client: %v", conn.RemoteAddr()) + i.logger.Info(fmt.Sprintf("agent.ipc: Accepted client: %v", conn.RemoteAddr())) metrics.IncrCounterWithLabels([]string{"agent", "ipc", "accept"}, 1, nil) // Wrap the connection in a client @@ -445,7 +447,7 @@ func (i *AgentIPC) handleClient(client *IPCClient) { // errors from Windows which appear to happen every // time there is an EOF. if err != io.EOF && !strings.Contains(strings.ToLower(err.Error()), "wsarecv") { - i.logger.Printf("[ERR] agent.ipc: failed to decode request header: %v", err) + i.logger.Error(fmt.Sprintf("agent.ipc: failed to decode request header: %v", err)) } } return @@ -453,7 +455,7 @@ func (i *AgentIPC) handleClient(client *IPCClient) { // Evaluate the command if err := i.handleRequest(client, &reqHeader); err != nil { - i.logger.Printf("[ERR] agent.ipc: Failed to evaluate request: %v", err) + i.logger.Error(fmt.Sprintf("agent.ipc: Failed to evaluate request: %v", err)) return } } @@ -475,7 +477,7 @@ func (i *AgentIPC) handleRequest(client *IPCClient, reqHeader *requestHeader) er // Ensure the client has authenticated after the handshake if necessary if i.authKey != "" && !client.didAuth && command != authCommand && command != handshakeCommand { - i.logger.Printf("[WARN] agent.ipc: Client sending commands before auth") + i.logger.Warn("agent.ipc: Client sending commands before auth") respHeader := responseHeader{Seq: seq, Error: authRequired} client.Send(&respHeader, nil) return nil @@ -931,12 +933,12 @@ func (i *AgentIPC) handleStop(client *IPCClient, seq uint64) error { } func (i *AgentIPC) handleLeave(client *IPCClient, seq uint64) error { - i.logger.Printf("[INFO] agent.ipc: Graceful leave triggered") + i.logger.Info("agent.ipc: Graceful leave triggered") // Do the leave err := i.agent.Leave() if err != nil { - i.logger.Printf("[ERR] agent.ipc: leave failed: %v", err) + i.logger.Error(fmt.Sprintf("agent.ipc: leave failed: %v", err)) } resp := responseHeader{Seq: seq, Error: errToString(err)} @@ -945,7 +947,7 @@ func (i *AgentIPC) handleLeave(client *IPCClient, seq uint64) error { // Trigger a shutdown! if err := i.agent.Shutdown(); err != nil { - i.logger.Printf("[ERR] agent.ipc: shutdown failed: %v", err) + i.logger.Error(fmt.Sprintf("agent.ipc: shutdown failed: %v", err)) } return err } diff --git a/cmd/serf/command/agent/ipc_event_stream.go b/cmd/serf/command/agent/ipc_event_stream.go index ad18c4e32..5b2b5e110 100644 --- a/cmd/serf/command/agent/ipc_event_stream.go +++ b/cmd/serf/command/agent/ipc_event_stream.go @@ -5,8 +5,8 @@ package agent import ( "fmt" - "log" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/serf/serf" ) @@ -20,11 +20,11 @@ type eventStream struct { client streamClient eventCh chan serf.Event filters []EventFilter - logger *log.Logger + logger hclog.Logger seq uint64 } -func newEventStream(client streamClient, filters []EventFilter, seq uint64, logger *log.Logger) *eventStream { +func newEventStream(client streamClient, filters []EventFilter, seq uint64, logger hclog.Logger) *eventStream { es := &eventStream{ client: client, eventCh: make(chan serf.Event, 512), @@ -50,7 +50,7 @@ HANDLE: select { case es.eventCh <- e: default: - es.logger.Printf("[WARN] agent.ipc: Dropping event to %v", es.client) + es.logger.Warn(fmt.Sprintf("agent.ipc: Dropping event to %v", es.client)) } } @@ -72,8 +72,8 @@ func (es *eventStream) stream() { err = fmt.Errorf("Unknown event type: %s", event.EventType().String()) } if err != nil { - es.logger.Printf("[ERR] agent.ipc: Failed to stream event to %v: %v", - es.client, err) + es.logger.Error(fmt.Sprintf("agent.ipc: Failed to stream event to %v: %v", + es.client, err)) return } } diff --git a/cmd/serf/command/agent/ipc_event_stream_test.go b/cmd/serf/command/agent/ipc_event_stream_test.go index 6def913ed..11eae2c1d 100644 --- a/cmd/serf/command/agent/ipc_event_stream_test.go +++ b/cmd/serf/command/agent/ipc_event_stream_test.go @@ -5,12 +5,12 @@ package agent import ( "bytes" - "log" "net" "os" "testing" "time" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/serf/serf" ) @@ -33,7 +33,9 @@ func (m *MockStreamClient) RegisterQuery(q *serf.Query) uint64 { func TestIPCEventStream(t *testing.T) { sc := &MockStreamClient{} filters := ParseEventFilter("user:foobar,member-join,query:deploy") - es := newEventStream(sc, filters, 42, log.New(os.Stderr, "", log.LstdFlags)) + es := newEventStream(sc, filters, 42, hclog.New(&hclog.LoggerOptions{ + Output: os.Stderr, + })) defer es.Stop() es.HandleEvent(serf.UserEvent{ diff --git a/cmd/serf/command/agent/ipc_log_stream.go b/cmd/serf/command/agent/ipc_log_stream.go index ce32bfaee..2e3c488ae 100644 --- a/cmd/serf/command/agent/ipc_log_stream.go +++ b/cmd/serf/command/agent/ipc_log_stream.go @@ -4,8 +4,9 @@ package agent import ( - "log" + "fmt" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/logutils" ) @@ -14,12 +15,12 @@ type logStream struct { client streamClient filter *logutils.LevelFilter logCh chan string - logger *log.Logger + logger hclog.Logger seq uint64 } func newLogStream(client streamClient, filter *logutils.LevelFilter, - seq uint64, logger *log.Logger) *logStream { + seq uint64, logger hclog.Logger) *logStream { ls := &logStream{ client: client, filter: filter, @@ -45,7 +46,7 @@ func (ls *logStream) HandleLog(l string) { // from the logWriter, and a log will need to invoke Write() which // already holds the lock. We must therefor do the log async, so // as to not deadlock - go ls.logger.Printf("[WARN] agent.ipc: Dropping logs to %v", ls.client) + go ls.logger.Warn(fmt.Sprintf("agent.ipc: Dropping logs to %v", ls.client)) } } @@ -60,8 +61,8 @@ func (ls *logStream) stream() { for line := range ls.logCh { rec.Log = line if err := ls.client.Send(&header, &rec); err != nil { - ls.logger.Printf("[ERR] agent.ipc: Failed to stream log to %v: %v", - ls.client, err) + ls.logger.Error(fmt.Sprintf("agent.ipc: Failed to stream log to %v: %v", + ls.client, err)) return } } diff --git a/cmd/serf/command/agent/ipc_log_stream_test.go b/cmd/serf/command/agent/ipc_log_stream_test.go index 47b6b0830..d5347d9a5 100644 --- a/cmd/serf/command/agent/ipc_log_stream_test.go +++ b/cmd/serf/command/agent/ipc_log_stream_test.go @@ -4,11 +4,11 @@ package agent import ( - "log" "os" "testing" "time" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/logutils" ) @@ -17,7 +17,9 @@ func TestIPCLogStream(t *testing.T) { filter := LevelFilter() filter.MinLevel = logutils.LogLevel("INFO") - ls := newLogStream(sc, filter, 42, log.New(os.Stderr, "", log.LstdFlags)) + ls := newLogStream(sc, filter, 42, hclog.New(&hclog.LoggerOptions{ + Output: os.Stderr, + })) defer ls.Stop() log := "[DEBUG] this is a test log" diff --git a/cmd/serf/command/agent/ipc_query_response_stream.go b/cmd/serf/command/agent/ipc_query_response_stream.go index 5a062bfbf..9eb228160 100644 --- a/cmd/serf/command/agent/ipc_query_response_stream.go +++ b/cmd/serf/command/agent/ipc_query_response_stream.go @@ -4,20 +4,21 @@ package agent import ( - "log" + "fmt" "time" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/serf/serf" ) // queryResponseStream is used to stream the query results back to a client type queryResponseStream struct { client streamClient - logger *log.Logger + logger hclog.Logger seq uint64 } -func newQueryResponseStream(client streamClient, seq uint64, logger *log.Logger) *queryResponseStream { +func newQueryResponseStream(client streamClient, seq uint64, logger hclog.Logger) *queryResponseStream { qs := &queryResponseStream{ client: client, logger: logger, @@ -38,17 +39,17 @@ func (qs *queryResponseStream) Stream(resp *serf.QueryResponse) { select { case a := <-ackCh: if err := qs.sendAck(a); err != nil { - qs.logger.Printf("[ERR] agent.ipc: Failed to stream ack to %v: %v", qs.client, err) + qs.logger.Error(fmt.Sprintf("agent.ipc: Failed to stream ack to %v: %v", qs.client, err)) return } case r := <-respCh: if err := qs.sendResponse(r.From, r.Payload); err != nil { - qs.logger.Printf("[ERR] agent.ipc: Failed to stream response to %v: %v", qs.client, err) + qs.logger.Error(fmt.Sprintf("agent.ipc: Failed to stream response to %v: %v", qs.client, err)) return } case <-done: if err := qs.sendDone(); err != nil { - qs.logger.Printf("[ERR] agent.ipc: Failed to stream query end to %v: %v", qs.client, err) + qs.logger.Error(fmt.Sprintf("agent.ipc: Failed to stream query end to %v: %v", qs.client, err)) } return } diff --git a/cmd/serf/command/agent/mdns.go b/cmd/serf/command/agent/mdns.go index 6ee6f5f4d..ba21644b3 100644 --- a/cmd/serf/command/agent/mdns.go +++ b/cmd/serf/command/agent/mdns.go @@ -6,10 +6,10 @@ package agent import ( "fmt" "io" - "log" "net" "time" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/mdns" ) @@ -23,7 +23,7 @@ const ( type AgentMDNS struct { agent *Agent discover string - logger *log.Logger + logger hclog.Logger seen map[string]struct{} server *mdns.Server replay bool @@ -62,11 +62,13 @@ func NewAgentMDNS(agent *Agent, logOutput io.Writer, replay bool, m := &AgentMDNS{ agent: agent, discover: discover, - logger: log.New(logOutput, "", log.LstdFlags), - seen: make(map[string]struct{}), - server: server, - replay: replay, - iface: iface, + logger: hclog.New(&hclog.LoggerOptions{ + Output: logOutput, + }), + seen: make(map[string]struct{}), + server: server, + replay: replay, + iface: iface, } // Start the background workers @@ -101,10 +103,10 @@ func (m *AgentMDNS) run() { // Attempt the join n, err := m.agent.Join(join, m.replay) if err != nil { - m.logger.Printf("[ERR] agent.mdns: Failed to join: %v", err) + m.logger.Error("agent.mdns: Failed to join: %v", err) } if n > 0 { - m.logger.Printf("[INFO] agent.mdns: Joined %d hosts", n) + m.logger.Info("agent.mdns: Joined %d hosts", n) } // Mark all as seen @@ -128,7 +130,7 @@ func (m *AgentMDNS) poll(hosts chan *mdns.ServiceEntry) { Entries: hosts, } if err := mdns.Query(¶ms); err != nil { - m.logger.Printf("[ERR] agent.mdns: Failed to poll for new hosts: %v", err) + m.logger.Error("agent.mdns: Failed to poll for new hosts: %v", err) } } diff --git a/cmd/serf/command/agent/syslog.go b/cmd/serf/command/agent/syslog.go index 85dd6f588..d71a7ca2a 100644 --- a/cmd/serf/command/agent/syslog.go +++ b/cmd/serf/command/agent/syslog.go @@ -6,8 +6,7 @@ package agent import ( "bytes" - "github.com/hashicorp/go-syslog" - "github.com/hashicorp/logutils" + gsyslog "github.com/hashicorp/go-syslog" ) // levelPriority is used to map a log level to a @@ -25,17 +24,11 @@ var levelPriority = map[string]gsyslog.Priority{ // writing them to a Syslogger. Implements the io.Writer // interface. type SyslogWrapper struct { - l gsyslog.Syslogger - filt *logutils.LevelFilter + l gsyslog.Syslogger } // Write is used to implement io.Writer func (s *SyslogWrapper) Write(p []byte) (int, error) { - // Skip syslog if the log level doesn't apply - if !s.filt.Check(p) { - return 0, nil - } - // Extract log level var level string afterLevel := p diff --git a/cmd/serf/command/agent/syslog_test.go b/cmd/serf/command/agent/syslog_test.go index 16e771be5..72abde9da 100644 --- a/cmd/serf/command/agent/syslog_test.go +++ b/cmd/serf/command/agent/syslog_test.go @@ -7,7 +7,7 @@ import ( "runtime" "testing" - "github.com/hashicorp/go-syslog" + gsyslog "github.com/hashicorp/go-syslog" "github.com/hashicorp/logutils" ) @@ -23,7 +23,7 @@ func TestSyslogFilter(t *testing.T) { filt := LevelFilter() filt.MinLevel = logutils.LogLevel("INFO") - s := &SyslogWrapper{l, filt} + s := &SyslogWrapper{l} n, err := s.Write([]byte("[INFO] test")) if err != nil { t.Fatalf("err: %v", err) @@ -31,12 +31,4 @@ func TestSyslogFilter(t *testing.T) { if n == 0 { t.Fatalf("should have logged") } - - n, err = s.Write([]byte("[DEBUG] test")) - if err != nil { - t.Fatalf("err: %v", err) - } - if n != 0 { - t.Fatalf("should not have logged") - } } diff --git a/coordinate/config.go b/coordinate/config.go index b6526953e..02cce7e98 100644 --- a/coordinate/config.go +++ b/coordinate/config.go @@ -14,12 +14,17 @@ import ( // here: // // [1] Dabek, Frank, et al. "Vivaldi: A decentralized network coordinate system." -// ACM SIGCOMM Computer Communication Review. Vol. 34. No. 4. ACM, 2004. +// +// ACM SIGCOMM Computer Communication Review. Vol. 34. No. 4. ACM, 2004. +// // [2] Ledlie, Jonathan, Paul Gardner, and Margo I. Seltzer. "Network Coordinates -// in the Wild." NSDI. Vol. 7. 2007. +// +// in the Wild." NSDI. Vol. 7. 2007. +// // [3] Lee, Sanghwan, et al. "On suitability of Euclidean embedding for -// host-based network coordinate systems." Networking, IEEE/ACM Transactions -// on 18.1 (2010): 27-40. +// +// host-based network coordinate systems." Networking, IEEE/ACM Transactions +// on 18.1 (2010): 27-40. type Config struct { // The dimensionality of the coordinate system. As discussed in [2], more // dimensions improves the accuracy of the estimates up to a point. Per [2] diff --git a/go.mod b/go.mod index 458403775..e55a91bb3 100644 --- a/go.mod +++ b/go.mod @@ -6,7 +6,7 @@ require ( github.com/armon/circbuf v0.0.0-20150827004946-bbbad097214e github.com/armon/go-metrics v0.4.1 github.com/armon/go-radix v1.0.0 // indirect - github.com/fatih/color v1.9.0 // indirect + github.com/hashicorp/go-hclog v1.5.0 github.com/hashicorp/go-msgpack v0.5.3 github.com/hashicorp/go-multierror v1.1.0 // indirect github.com/hashicorp/go-syslog v1.0.0 @@ -14,7 +14,6 @@ require ( github.com/hashicorp/logutils v1.0.0 github.com/hashicorp/mdns v1.0.4 github.com/hashicorp/memberlist v0.5.0 - github.com/mattn/go-colorable v0.1.6 // indirect github.com/mitchellh/cli v1.1.5 github.com/mitchellh/mapstructure v1.5.0 github.com/posener/complete v1.2.3 // indirect diff --git a/go.sum b/go.sum index eb8f9ef7b..fb1544c77 100644 --- a/go.sum +++ b/go.sum @@ -29,8 +29,8 @@ github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSs github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/fatih/color v1.7.0/go.mod h1:Zm6kSWBoL9eyXnKyktHP6abPY2pDugNf5KwzbycvMj4= -github.com/fatih/color v1.9.0 h1:8xPHl4/q1VyqGIPif1F+1V3Y3lSmrq01EabUW3CoW5s= -github.com/fatih/color v1.9.0/go.mod h1:eQcE1qtQxscV5RaZvpXrrb8Drkc3/DdQ+uUYCNjL+zU= +github.com/fatih/color v1.13.0 h1:8LOYc1KYPPmyKMuN8QV2DNRWNbLo6LZ0iLs8+mlH53w= +github.com/fatih/color v1.13.0/go.mod h1:kLAiJbzzSOZDVNGyDpeOxJ47H46qBXwg5ILebYFFOfk= github.com/go-kit/kit v0.8.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= github.com/go-kit/kit v0.9.0/go.mod h1:xBxKIO96dXMWWy0MnWVtmwkA9/13aqxPnvrjFYMA2as= github.com/go-logfmt/logfmt v0.3.0/go.mod h1:Qt1PoO58o5twSAckw1HlFXLmHsOX5/0LbT9GBnD5lWE= @@ -51,6 +51,8 @@ github.com/google/uuid v1.1.2/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+ github.com/hashicorp/errwrap v1.0.0 h1:hLrqtEDnRye3+sgx6z4qVLNuviH3MR5aQ0ykNJa/UYA= github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brvKWEqk/Jc4= github.com/hashicorp/go-cleanhttp v0.5.0/go.mod h1:JpRdi6/HCYpAwUzNwuwqhbovhLtngrth3wmdIIUrZ80= +github.com/hashicorp/go-hclog v1.5.0 h1:bI2ocEMgcVlz55Oj1xZNBsVi900c7II+fWDyV9o+13c= +github.com/hashicorp/go-hclog v1.5.0/go.mod h1:W4Qnvbt70Wk/zYJryRzDRU/4r0kIg0PVHBcfoyhpF5M= github.com/hashicorp/go-immutable-radix v1.0.0 h1:AKDB1HM5PWEA7i4nhcpwOrO2byshxBjXVn/J/3+z5/0= github.com/hashicorp/go-immutable-radix v1.0.0/go.mod h1:0y9vanUI8NX6FsYoO3zeMjhV/C5i9g4Q3DwcSNZ4P60= github.com/hashicorp/go-msgpack v0.5.3 h1:zKjpN5BK/P5lMYrLmBHdBULWbJ0XpYR+7NGzqkZzoD4= @@ -90,14 +92,13 @@ github.com/kr/pty v1.1.1/go.mod h1:pFQYn66WHrOpPYNljwOMqo10TkYh1fy3cYio2l3bCsQ= github.com/kr/text v0.1.0 h1:45sCR5RtlFHMR4UwH9sdQ5TC8v0qDQCHnXt+kaKSTVE= github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI= github.com/mattn/go-colorable v0.0.9/go.mod h1:9vuHe8Xs5qXnSaW/c/ABM9alt+Vo+STaOChaDxuIBZU= -github.com/mattn/go-colorable v0.1.4/go.mod h1:U0ppj6V5qS13XJ6of8GYAs25YV2eR4EVcfRqFIhoBtE= -github.com/mattn/go-colorable v0.1.6 h1:6Su7aK7lXmJ/U79bYtBjLNaha4Fs1Rg9plHpcH+vvnE= -github.com/mattn/go-colorable v0.1.6/go.mod h1:u6P/XSegPjTcexA+o6vUJrdnUu04hMope9wVRipJSqc= +github.com/mattn/go-colorable v0.1.9/go.mod h1:u6P/XSegPjTcexA+o6vUJrdnUu04hMope9wVRipJSqc= +github.com/mattn/go-colorable v0.1.12 h1:jF+Du6AlPIjs2BiUiQlKOX0rt3SujHxPnksPKZbaA40= +github.com/mattn/go-colorable v0.1.12/go.mod h1:u5H1YNBxpqRaxsYJYSkiCWKzEfiAb1Gb520KVy5xxl4= github.com/mattn/go-isatty v0.0.3/go.mod h1:M+lRXTBqGeGNdLjl/ufCoiOlB5xdOkqRJdNxMWT7Zi4= -github.com/mattn/go-isatty v0.0.8/go.mod h1:Iq45c/XA43vh69/j3iqttzPXn0bhXyGjM0Hdxcsrc5s= -github.com/mattn/go-isatty v0.0.11/go.mod h1:PhnuNfih5lzO57/f3n+odYbM4JtupLOxQOAqxQCu2WE= -github.com/mattn/go-isatty v0.0.12 h1:wuysRhFDzyxgEmMf5xjvJ2M9dZoWAXNNr5LSBS7uHXY= github.com/mattn/go-isatty v0.0.12/go.mod h1:cbi8OIDigv2wuxKPP5vlRcQ1OAZbq2CE4Kysco4FUpU= +github.com/mattn/go-isatty v0.0.14 h1:yVuAays6BHfxijgZPzw+3Zlu5yQgKGP2/hcQbHb7S9Y= +github.com/mattn/go-isatty v0.0.14/go.mod h1:7GGIvUiUoEMVVmxf/4nioHXj79iQHKdU27kJ6hsGG94= github.com/matttproud/golang_protobuf_extensions v1.0.1/go.mod h1:D8He9yQNgCq6Z5Ld7szi9bcBfOoFv/3dc6xSMkL2PC0= github.com/miekg/dns v1.1.26/go.mod h1:bPDLeHnStXmXAq1m/Ch/hvfNHr14JKNPMBo3VZKjuso= github.com/miekg/dns v1.1.41 h1:WMszZWJG0XmzbK9FEmzH2TVcqYzFesusSIB41b8KHxY= @@ -152,8 +153,9 @@ github.com/stretchr/testify v1.2.2/go.mod h1:a8OnRcib4nhh0OaRAV+Yts87kKdq0PP7pXf github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4= github.com/stretchr/testify v1.5.1/go.mod h1:5W2xD1RspED5o8YsWQXVCued0rvSQ+mT+I5cxcmMvtA= -github.com/stretchr/testify v1.6.1 h1:hDPOHmpOpP40lSULcqw7IrRb/u7w6RpDC9399XyoNd0= github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.7.2 h1:4jaiDzPyXQvSd7D0EjG45355tLlV3VOECpq10pLC+8s= +github.com/stretchr/testify v1.7.2/go.mod h1:R6va5+xMeoiuVRoj+gSkQ7d3FALtqAAGI1FQKckRals= github.com/tv42/httpunix v0.0.0-20150427012821-b75d8614f926/go.mod h1:9ESjWnEqriFuLhtthL60Sar/7RFoluCcXsuvEwTV5KM= golang.org/x/crypto v0.0.0-20180904163835-0709b304e793/go.mod h1:6SG95UA2DQfeDnfUPMdvaQW0Q7yPrPDi9nlGo2tz2b4= golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w= @@ -178,18 +180,19 @@ golang.org/x/sync v0.0.0-20210220032951-036812b2e83c/go.mod h1:RxMgew5VJxzue5/jJ golang.org/x/sys v0.0.0-20180905080454-ebe1bf3edb33/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20181116152217-5ac8a444bdc5/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20190215142949-d0b11bdaac8a/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= -golang.org/x/sys v0.0.0-20190222072716-a9d3bda3a223/go.mod h1:STP8DvDyc/dI5b8T5hshtkjS+E42TnysNCUPdjciGhY= golang.org/x/sys v0.0.0-20190412213103-97732733099d/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190422165155-953cdadca894/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190922100055-0a153f010e69/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20190924154521-2837fb4f24fe/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= -golang.org/x/sys v0.0.0-20191026070338-33540a1f6037/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200116001909-b77594299b42/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200122134326-e047566fdf82/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20200223170610-d5e6a3e2c0ae/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20201119102817-f84b799fce68/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210303074136-134d130e1a04/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= golang.org/x/sys v0.0.0-20210330210617-4fbd30eecc44/go.mod h1:h1NjWce9XRLGQEsW7wpKNCjG9DtNlClVuFLEZdDNbEs= +golang.org/x/sys v0.0.0-20210630005230-0f9fa26af87c/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20210927094055-39ccf1dd6fa6/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= +golang.org/x/sys v0.0.0-20220503163025-988cb79eb6c6/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/sys v0.0.0-20220728004956-3c1f35247d10 h1:WIoqL4EROvwiPdUtaip4VcDdpZ4kha7wBWZrbVKCIZg= golang.org/x/sys v0.0.0-20220728004956-3c1f35247d10/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg= golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo= @@ -211,5 +214,6 @@ gopkg.in/yaml.v2 v2.2.4/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= gopkg.in/yaml.v2 v2.2.5/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= gopkg.in/yaml.v2 v2.3.0 h1:clyUAQHOM3G0M3f5vQj7LuJrETvjVot3Z5el9nffUtU= gopkg.in/yaml.v2 v2.3.0/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI= -gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c h1:dUUwHk2QECo/6vqA44rthZ8ie2QXMNeKRTHCNY2nXvo= gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/serf/config.go b/serf/config.go index 2a4fb7649..fc38bdc92 100644 --- a/serf/config.go +++ b/serf/config.go @@ -5,11 +5,11 @@ package serf import ( "io" - "log" "os" "time" "github.com/armon/go-metrics" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/memberlist" ) @@ -209,7 +209,7 @@ type Config struct { // this for the internal logger. If Logger is not set, it will fall back to the // behavior for using LogOutput. You cannot specify both LogOutput and Logger // at the same time. - Logger *log.Logger + Logger hclog.Logger // SnapshotPath if provided is used to snapshot live nodes as well // as lamport clock values. When Serf is started with a snapshot, diff --git a/serf/delegate.go b/serf/delegate.go index 6555f5a26..6db6b36f1 100644 --- a/serf/delegate.go +++ b/serf/delegate.go @@ -47,53 +47,53 @@ func (d *delegate) NotifyMsg(buf []byte) { case messageLeaveType: var leave messageLeave if err := decodeMessage(buf[1:], &leave); err != nil { - d.serf.logger.Printf("[ERR] serf: Error decoding leave message: %s", err) + d.serf.logger.Error(fmt.Sprintf("serf: Error decoding leave message: %s", err)) break } - d.serf.logger.Printf("[DEBUG] serf: messageLeaveType: %s", leave.Node) + d.serf.logger.Debug(fmt.Sprintf("serf: messageLeaveType: %s", leave.Node)) rebroadcast = d.serf.handleNodeLeaveIntent(&leave) case messageJoinType: var join messageJoin if err := decodeMessage(buf[1:], &join); err != nil { - d.serf.logger.Printf("[ERR] serf: Error decoding join message: %s", err) + d.serf.logger.Error(fmt.Sprintf("serf: Error decoding join message: %s", err)) break } - d.serf.logger.Printf("[DEBUG] serf: messageJoinType: %s", join.Node) + d.serf.logger.Debug(fmt.Sprintf("serf: messageJoinType: %s", join.Node)) rebroadcast = d.serf.handleNodeJoinIntent(&join) case messageUserEventType: var event messageUserEvent if err := decodeMessage(buf[1:], &event); err != nil { - d.serf.logger.Printf("[ERR] serf: Error decoding user event message: %s", err) + d.serf.logger.Error(fmt.Sprintf("serf: Error decoding user event message: %s", err)) break } - d.serf.logger.Printf("[DEBUG] serf: messageUserEventType: %s", event.Name) + d.serf.logger.Debug(fmt.Sprintf("serf: messageUserEventType: %s", event.Name)) rebroadcast = d.serf.handleUserEvent(&event) rebroadcastQueue = d.serf.eventBroadcasts case messageQueryType: var query messageQuery if err := decodeMessage(buf[1:], &query); err != nil { - d.serf.logger.Printf("[ERR] serf: Error decoding query message: %s", err) + d.serf.logger.Error("serf: Error decoding query message: %s", err) break } - d.serf.logger.Printf("[DEBUG] serf: messageQueryType: %s", query.Name) + d.serf.logger.Debug(fmt.Sprintf("serf: messageQueryType: %s", query.Name)) rebroadcast = d.serf.handleQuery(&query) rebroadcastQueue = d.serf.queryBroadcasts case messageQueryResponseType: var resp messageQueryResponse if err := decodeMessage(buf[1:], &resp); err != nil { - d.serf.logger.Printf("[ERR] serf: Error decoding query response message: %s", err) + d.serf.logger.Error(fmt.Sprintf("serf: Error decoding query response message: %s", err)) break } - d.serf.logger.Printf("[DEBUG] serf: messageQueryResponseType: %v", resp.From) + d.serf.logger.Debug(fmt.Sprintf("serf: messageQueryResponseType: %v", resp.From)) d.serf.handleQueryResponse(&resp) case messageRelayType: @@ -102,7 +102,7 @@ func (d *delegate) NotifyMsg(buf []byte) { reader := bytes.NewReader(buf[1:]) decoder := codec.NewDecoder(reader, &handle) if err := decoder.Decode(&header); err != nil { - d.serf.logger.Printf("[ERR] serf: Error decoding relay header: %s", err) + d.serf.logger.Error(fmt.Sprintf("serf: Error decoding relay header: %s", err)) break } @@ -115,14 +115,14 @@ func (d *delegate) NotifyMsg(buf []byte) { Name: header.DestName, } - d.serf.logger.Printf("[DEBUG] serf: Relaying response to addr: %s", header.DestAddr.String()) + d.serf.logger.Debug(fmt.Sprintf("serf: Relaying response to addr: %s", header.DestAddr.String())) if err := d.serf.memberlist.SendToAddress(addr, raw); err != nil { - d.serf.logger.Printf("[ERR] serf: Error forwarding message to %s: %s", header.DestAddr.String(), err) + d.serf.logger.Error(fmt.Sprintf("serf: Error forwarding message to %s: %s", header.DestAddr.String(), err)) break } default: - d.serf.logger.Printf("[WARN] serf: Received message of unknown type: %d", t) + d.serf.logger.Warn(fmt.Sprintf("serf: Received message of unknown type: %d", t)) } if rebroadcast { @@ -202,7 +202,7 @@ func (d *delegate) LocalState(join bool) []byte { // Encode the push pull state buf, err := encodeMessage(messagePushPullType, &pp) if err != nil { - d.serf.logger.Printf("[ERR] serf: Failed to encode local state: %v", err) + d.serf.logger.Error("serf: Failed to encode local state: %v", err) return nil } return buf @@ -211,13 +211,13 @@ func (d *delegate) LocalState(join bool) []byte { func (d *delegate) MergeRemoteState(buf []byte, isJoin bool) { // Ensure we have a message if len(buf) == 0 { - d.serf.logger.Printf("[ERR] serf: Remote state is zero bytes") + d.serf.logger.Error("serf: Remote state is zero bytes") return } // Check the message type if messageType(buf[0]) != messagePushPullType { - d.serf.logger.Printf("[ERR] serf: Remote state has bad type prefix: %v", buf[0]) + d.serf.logger.Error("serf: Remote state has bad type prefix: %v", buf[0]) return } @@ -228,7 +228,7 @@ func (d *delegate) MergeRemoteState(buf []byte, isJoin bool) { // Attempt a decode pp := messagePushPull{} if err := decodeMessage(buf[1:], &pp); err != nil { - d.serf.logger.Printf("[ERR] serf: Failed to decode remote state: %v", err) + d.serf.logger.Error("serf: Failed to decode remote state: %v", err) return } diff --git a/serf/internal_query.go b/serf/internal_query.go index cac3340cf..4b433ab66 100644 --- a/serf/internal_query.go +++ b/serf/internal_query.go @@ -6,8 +6,9 @@ package serf import ( "encoding/base64" "fmt" - "log" "strings" + + "github.com/hashicorp/go-hclog" ) const ( @@ -50,7 +51,7 @@ func internalQueryName(name string) string { // _serf and respond to them as appropriate. type serfQueries struct { inCh chan Event - logger *log.Logger + logger hclog.Logger outCh chan<- Event serf *Serf shutdownCh <-chan struct{} @@ -75,7 +76,7 @@ type nodeKeyResponse struct { // newSerfQueries is used to create a new serfQueries. We return an event // channel that is ingested and forwarded to an outCh. Any Queries that // have the InternalQueryPrefix are handled instead of forwarded. -func newSerfQueries(serf *Serf, logger *log.Logger, outCh chan<- Event, shutdownCh <-chan struct{}) (chan<- Event, error) { +func newSerfQueries(serf *Serf, logger hclog.Logger, outCh chan<- Event, shutdownCh <-chan struct{}) (chan<- Event, error) { inCh := make(chan Event, 1024) q := &serfQueries{ inCh: inCh, @@ -125,7 +126,7 @@ func (s *serfQueries) handleQuery(q *Query) { case listKeysQuery: s.handleListKeys(q) default: - s.logger.Printf("[WARN] serf: Unhandled internal query '%s'", queryName) + s.logger.Warn(fmt.Sprintf("serf: Unhandled internal query '%s'", queryName)) } } @@ -140,7 +141,7 @@ func (s *serfQueries) handleConflict(q *Query) { if node == s.serf.config.NodeName { return } - s.logger.Printf("[DEBUG] serf: Got conflict resolution query for '%s'", node) + s.logger.Debug(fmt.Sprintf("serf: Got conflict resolution query for '%s'", node)) // Look for the member info var out *Member @@ -153,13 +154,13 @@ func (s *serfQueries) handleConflict(q *Query) { // Encode the response buf, err := encodeMessage(messageConflictResponseType, out) if err != nil { - s.logger.Printf("[ERR] serf: Failed to encode conflict query response: %v", err) + s.logger.Error("serf: Failed to encode conflict query response: %v", err) return } // Send our answer if err := q.Respond(buf); err != nil { - s.logger.Printf("[ERR] serf: Failed to respond to conflict query: %v", err) + s.logger.Error("serf: Failed to respond to conflict query: %v", err) } } @@ -196,7 +197,7 @@ func (s *serfQueries) keyListResponseWithCorrectSize(q *Query, resp *nodeKeyResp } if actual > i { - s.logger.Printf("[WARN] serf: %s", resp.Message) + s.logger.Warn(fmt.Sprintf("serf: %s", resp.Message)) } return raw, qresp, nil } @@ -209,21 +210,21 @@ func (s *serfQueries) sendKeyResponse(q *Query, resp *nodeKeyResponse) { case internalQueryName(listKeysQuery): raw, qresp, err := s.keyListResponseWithCorrectSize(q, resp) if err != nil { - s.logger.Printf("[ERR] serf: %v", err) + s.logger.Error(fmt.Sprintf("serf: %v", err)) return } if err := q.respondWithMessageAndResponse(raw, qresp); err != nil { - s.logger.Printf("[ERR] serf: Failed to respond to key query: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to respond to key query: %v", err)) return } default: buf, err := encodeMessage(messageKeyResponseType, resp) if err != nil { - s.logger.Printf("[ERR] serf: Failed to encode key response: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to encode key response: %v", err)) return } if err := q.Respond(buf); err != nil { - s.logger.Printf("[ERR] serf: Failed to respond to key query: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to respond to key query: %v", err)) return } } @@ -241,27 +242,27 @@ func (s *serfQueries) handleInstallKey(q *Query) { err := decodeMessage(q.Payload[1:], &req) if err != nil { - s.logger.Printf("[ERR] serf: Failed to decode key request: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to decode key request: %v", err)) goto SEND } if !s.serf.EncryptionEnabled() { response.Message = "No keyring to modify (encryption not enabled)" - s.logger.Printf("[ERR] serf: No keyring to modify (encryption not enabled)") + s.logger.Error(fmt.Sprintf("serf: No keyring to modify (encryption not enabled)")) goto SEND } - s.logger.Printf("[INFO] serf: Received install-key query") + s.logger.Info("serf: Received install-key query") if err := keyring.AddKey(req.Key); err != nil { response.Message = err.Error() - s.logger.Printf("[ERR] serf: Failed to install key: %s", err) + s.logger.Error(fmt.Sprintf("serf: Failed to install key: %s", err)) goto SEND } if s.serf.config.KeyringFile != "" { if err := s.serf.writeKeyringFile(); err != nil { response.Message = err.Error() - s.logger.Printf("[ERR] serf: Failed to write keyring file: %s", err) + s.logger.Error(fmt.Sprintf("serf: Failed to write keyring file: %s", err)) goto SEND } } @@ -283,26 +284,26 @@ func (s *serfQueries) handleUseKey(q *Query) { err := decodeMessage(q.Payload[1:], &req) if err != nil { - s.logger.Printf("[ERR] serf: Failed to decode key request: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to decode key request: %v", err)) goto SEND } if !s.serf.EncryptionEnabled() { response.Message = "No keyring to modify (encryption not enabled)" - s.logger.Printf("[ERR] serf: No keyring to modify (encryption not enabled)") + s.logger.Error(fmt.Sprintf("serf: No keyring to modify (encryption not enabled)")) goto SEND } - s.logger.Printf("[INFO] serf: Received use-key query") + s.logger.Info("serf: Received use-key query") if err := keyring.UseKey(req.Key); err != nil { response.Message = err.Error() - s.logger.Printf("[ERR] serf: Failed to change primary key: %s", err) + s.logger.Error(fmt.Sprintf("serf: Failed to change primary key: %s", err)) goto SEND } if err := s.serf.writeKeyringFile(); err != nil { response.Message = err.Error() - s.logger.Printf("[ERR] serf: Failed to write keyring file: %s", err) + s.logger.Error(fmt.Sprintf("serf: Failed to write keyring file: %s", err)) goto SEND } @@ -323,26 +324,26 @@ func (s *serfQueries) handleRemoveKey(q *Query) { err := decodeMessage(q.Payload[1:], &req) if err != nil { - s.logger.Printf("[ERR] serf: Failed to decode key request: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to decode key request: %v", err)) goto SEND } if !s.serf.EncryptionEnabled() { response.Message = "No keyring to modify (encryption not enabled)" - s.logger.Printf("[ERR] serf: No keyring to modify (encryption not enabled)") + s.logger.Error(fmt.Sprintf("serf: No keyring to modify (encryption not enabled)")) goto SEND } - s.logger.Printf("[INFO] serf: Received remove-key query") + s.logger.Info("serf: Received remove-key query") if err := keyring.RemoveKey(req.Key); err != nil { response.Message = err.Error() - s.logger.Printf("[ERR] serf: Failed to remove key: %s", err) + s.logger.Error(fmt.Sprintf("serf: Failed to remove key: %s", err)) goto SEND } if err := s.serf.writeKeyringFile(); err != nil { response.Message = err.Error() - s.logger.Printf("[ERR] serf: Failed to write keyring file: %s", err) + s.logger.Error(fmt.Sprintf("serf: Failed to write keyring file: %s", err)) goto SEND } @@ -362,11 +363,11 @@ func (s *serfQueries) handleListKeys(q *Query) { var primaryKeyBytes []byte if !s.serf.EncryptionEnabled() { response.Message = "Keyring is empty (encryption not enabled)" - s.logger.Printf("[ERR] serf: Keyring is empty (encryption not enabled)") + s.logger.Error("serf: Keyring is empty (encryption not enabled)") goto SEND } - s.logger.Printf("[INFO] serf: Received list-keys query") + s.logger.Info("serf: Received list-keys query") for _, keyBytes := range keyring.GetKeys() { // Encode the keys before sending the response. This should help take // some the burden of doing this off of the asking member. diff --git a/serf/internal_query_test.go b/serf/internal_query_test.go index 470e9e7f5..b449215b9 100644 --- a/serf/internal_query_test.go +++ b/serf/internal_query_test.go @@ -4,11 +4,11 @@ package serf import ( - "log" - "os" "strings" "testing" "time" + + "github.com/hashicorp/go-hclog" ) func TestInternalQueryName(t *testing.T) { @@ -20,7 +20,7 @@ func TestInternalQueryName(t *testing.T) { func TestSerfQueries_Passthrough(t *testing.T) { serf := &Serf{} - logger := log.New(os.Stderr, "", log.LstdFlags) + logger := hclog.Default() outCh := make(chan Event, 4) shutdown := make(chan struct{}) defer close(shutdown) @@ -50,7 +50,7 @@ func TestSerfQueries_Passthrough(t *testing.T) { func TestSerfQueries_Ping(t *testing.T) { serf := &Serf{} - logger := log.New(os.Stderr, "", log.LstdFlags) + logger := hclog.Default() outCh := make(chan Event, 4) shutdown := make(chan struct{}) defer close(shutdown) @@ -72,7 +72,7 @@ func TestSerfQueries_Ping(t *testing.T) { func TestSerfQueries_Conflict_SameName(t *testing.T) { serf := &Serf{config: &Config{NodeName: "foo"}} - logger := log.New(os.Stderr, "", log.LstdFlags) + logger := hclog.Default() outCh := make(chan Event, 4) shutdown := make(chan struct{}) defer close(shutdown) @@ -124,7 +124,7 @@ func TestSerfQueries_estimateMaxKeysInListKeyResponseFactor(t *testing.T) { } func TestSerfQueries_keyListResponseWithCorrectSize(t *testing.T) { - s := serfQueries{logger: log.New(os.Stderr, "", log.LstdFlags)} + s := serfQueries{logger: hclog.Default()} q := Query{id: 0, serf: &Serf{config: &Config{NodeName: "", QueryResponseSizeLimit: 1024}}} cases := []struct { resp nodeKeyResponse diff --git a/serf/keymanager.go b/serf/keymanager.go index 5ce8fbf8d..0ce370805 100644 --- a/serf/keymanager.go +++ b/serf/keymanager.go @@ -77,7 +77,7 @@ func (k *KeyManager) streamKeyResp(resp *KeyResponse, ch <-chan NodeResponse) { if nodeResponse.Result && len(nodeResponse.Message) > 0 { resp.Messages[r.From] = nodeResponse.Message - k.serf.logger.Println("[WARN] serf:", nodeResponse.Message) + k.serf.logger.Warn(fmt.Sprintf("serf: %s", nodeResponse.Message)) } // Currently only used for key list queries, this adds keys to a counter diff --git a/serf/ping_delegate.go b/serf/ping_delegate.go index 78a34f742..6c42113d0 100644 --- a/serf/ping_delegate.go +++ b/serf/ping_delegate.go @@ -5,6 +5,7 @@ package serf import ( "bytes" + "fmt" "time" "github.com/armon/go-metrics" @@ -39,7 +40,7 @@ func (p *pingDelegate) AckPayload() []byte { // The rest of the message is the serialized coordinate. enc := codec.NewEncoder(&buf, &codec.MsgpackHandle{}) if err := enc.Encode(p.serf.coordClient.GetCoordinate()); err != nil { - p.serf.logger.Printf("[ERR] serf: Failed to encode coordinate: %v\n", err) + p.serf.logger.Error(fmt.Sprintf("serf: Failed to encode coordinate: %v\n", err)) } return buf.Bytes() } @@ -54,7 +55,7 @@ func (p *pingDelegate) NotifyPingComplete(other *memberlist.Node, rtt time.Durat // Verify ping version in the header. version := payload[0] if version != PingVersion { - p.serf.logger.Printf("[ERR] serf: Unsupported ping version: %v", version) + p.serf.logger.Error(fmt.Sprintf("serf: Unsupported ping version: %v", version)) return } @@ -63,7 +64,7 @@ func (p *pingDelegate) NotifyPingComplete(other *memberlist.Node, rtt time.Durat dec := codec.NewDecoder(r, &codec.MsgpackHandle{}) var coord coordinate.Coordinate if err := dec.Decode(&coord); err != nil { - p.serf.logger.Printf("[ERR] serf: Failed to decode coordinate from ping: %v", err) + p.serf.logger.Error(fmt.Sprintf("serf: Failed to decode coordinate from ping: %v", err)) return } @@ -72,8 +73,8 @@ func (p *pingDelegate) NotifyPingComplete(other *memberlist.Node, rtt time.Durat after, err := p.serf.coordClient.Update(other.Name, &coord, rtt) if err != nil { metrics.IncrCounterWithLabels([]string{"serf", "coordinate", "rejected"}, 1, p.serf.metricLabels) - p.serf.logger.Printf("[TRACE] serf: Rejected coordinate from %s: %v\n", - other.Name, err) + p.serf.logger.Trace(fmt.Sprintf("serf: Rejected coordinate from %s: %v\n", + other.Name, err)) return } diff --git a/serf/query.go b/serf/query.go index 2c04a9078..6e713b94c 100644 --- a/serf/query.go +++ b/serf/query.go @@ -219,7 +219,7 @@ func (s *Serf) shouldProcessQuery(filters [][]byte) bool { // Decode the filter var nodes filterNode if err := decodeMessage(filter[1:], &nodes); err != nil { - s.logger.Printf("[WARN] serf: failed to decode filterNodeType: %v", err) + s.logger.Warn(fmt.Sprintf("serf: failed to decode filterNodeType: %v", err)) return false } @@ -239,7 +239,7 @@ func (s *Serf) shouldProcessQuery(filters [][]byte) bool { // Decode the filter var filt filterTag if err := decodeMessage(filter[1:], &filt); err != nil { - s.logger.Printf("[WARN] serf: failed to decode filterTagType: %v", err) + s.logger.Warn(fmt.Sprintf("serf: failed to decode filterTagType: %v", err)) return false } @@ -247,7 +247,7 @@ func (s *Serf) shouldProcessQuery(filters [][]byte) bool { tags := s.config.Tags matched, err := regexp.MatchString(filt.Expr, tags[filt.Tag]) if err != nil { - s.logger.Printf("[WARN] serf: failed to compile filter regex (%s): %v", filt.Expr, err) + s.logger.Warn("serf: failed to compile filter regex (%s): %v", filt.Expr, err) return false } if !matched { @@ -255,7 +255,7 @@ func (s *Serf) shouldProcessQuery(filters [][]byte) bool { } default: - s.logger.Printf("[WARN] serf: query has unrecognized filter type: %d", filter[0]) + s.logger.Warn(fmt.Sprintf("serf: query has unrecognized filter type: %d", filter[0])) return false } } diff --git a/serf/serf.go b/serf/serf.go index 96f1d0a09..2b185fce2 100644 --- a/serf/serf.go +++ b/serf/serf.go @@ -9,7 +9,6 @@ import ( "encoding/json" "fmt" "io/ioutil" - "log" "math/rand" "net" "os" @@ -20,6 +19,7 @@ import ( "time" "github.com/armon/go-metrics" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/go-msgpack/codec" "github.com/hashicorp/memberlist" "github.com/hashicorp/serf/coordinate" @@ -96,7 +96,7 @@ type Serf struct { queryResponse map[LamportTime]*QueryResponse queryLock sync.RWMutex - logger *log.Logger + logger hclog.Logger joinLock sync.Mutex stateLock sync.Mutex state SerfState @@ -266,7 +266,10 @@ func Create(conf *Config) (*Serf, error) { if logOutput == nil { logOutput = os.Stderr } - logger = log.New(logOutput, "", log.LstdFlags) + + logger = hclog.New(&hclog.LoggerOptions{ + Output: logOutput, + }) } serf := &Serf{ @@ -688,7 +691,7 @@ func (s *Serf) broadcastJoin(ltime LamportTime) error { // Start broadcasting the update if err := s.broadcast(messageJoinType, &msg, nil); err != nil { - s.logger.Printf("[WARN] serf: Failed to broadcast join intent: %v", err) + s.logger.Warn(fmt.Sprintf("serf: Failed to broadcast join intent: %v", err)) return err } return nil @@ -739,14 +742,14 @@ func (s *Serf) Leave() error { select { case <-notifyCh: case <-time.After(s.config.BroadcastTimeout): - s.logger.Printf("[WARN] serf: timeout while waiting for graceful leave") + s.logger.Warn("serf: timeout while waiting for graceful leave") } } // Attempt the memberlist leave err := s.memberlist.Leave(s.config.BroadcastTimeout) if err != nil { - s.logger.Printf("[WARN] serf: timeout waiting for leave broadcast: %s", err.Error()) + s.logger.Warn(fmt.Sprintf("serf: timeout waiting for leave broadcast: %s", err.Error())) } // Wait for the leave to propagate through the cluster. The broadcast @@ -870,7 +873,7 @@ func (s *Serf) Shutdown() error { } if s.state != SerfLeft { - s.logger.Printf("[WARN] serf: Shutdown without a Leave") + s.logger.Warn("serf: Shutdown without a Leave") } // Wait to close the shutdown channel until after we've shut down the @@ -995,8 +998,8 @@ func (s *Serf) handleNodeJoin(n *memberlist.Node) { metrics.IncrCounterWithLabels([]string{"serf", "member", "join"}, 1, s.metricLabels) // Send an event along - s.logger.Printf("[INFO] serf: EventMemberJoin: %s %s", - member.Member.Name, member.Member.Addr) + s.logger.Info(fmt.Sprintf("serf: EventMemberJoin: %s %s", + member.Member.Name, member.Member.Addr)) if s.config.EventCh != nil { s.config.EventCh <- MemberEvent{ Type: EventMemberJoin, @@ -1029,7 +1032,7 @@ func (s *Serf) handleNodeLeave(n *memberlist.Node) { s.failedMembers = append(s.failedMembers, member) default: // Unknown state that it was in? Just don't do anything - s.logger.Printf("[WARN] serf: Bad state when leave: %d", member.Status) + s.logger.Warn(fmt.Sprintf("serf: Bad state when leave: %d", member.Status)) return } @@ -1044,8 +1047,8 @@ func (s *Serf) handleNodeLeave(n *memberlist.Node) { // Update some metrics metrics.IncrCounterWithLabels([]string{"serf", "member", member.Status.String()}, 1, s.metricLabels) - s.logger.Printf("[INFO] serf: %s: %s %s", - eventStr, member.Member.Name, member.Member.Addr) + s.logger.Info(fmt.Sprintf("serf: %s: %s %s", + eventStr, member.Member.Name, member.Member.Addr)) if s.config.EventCh != nil { s.config.EventCh <- MemberEvent{ Type: event, @@ -1089,7 +1092,7 @@ func (s *Serf) handleNodeUpdate(n *memberlist.Node) { metrics.IncrCounterWithLabels([]string{"serf", "member", "update"}, 1, s.metricLabels) // Send an event along - s.logger.Printf("[INFO] serf: EventMemberUpdate: %s", member.Member.Name) + s.logger.Info(fmt.Sprintf("serf: EventMemberUpdate: %s", member.Member.Name)) if s.config.EventCh != nil { s.config.EventCh <- MemberEvent{ Type: EventMemberUpdate, @@ -1122,7 +1125,7 @@ func (s *Serf) handleNodeLeaveIntent(leaveMsg *messageLeave) bool { // Refute us leaving if we are in the alive state // Must be done in another goroutine since we have the memberLock if leaveMsg.Node == s.config.NodeName && state == SerfAlive { - s.logger.Printf("[DEBUG] serf: Refuting an older leave intent") + s.logger.Debug("serf: Refuting an older leave intent") go s.broadcastJoin(s.clock.Time()) return false } @@ -1166,8 +1169,8 @@ func (s *Serf) handleNodeLeaveIntent(leaveMsg *messageLeave) bool { // We must push a message indicating the node has now // left to allow higher-level applications to handle the // graceful leave. - s.logger.Printf("[INFO] serf: EventMemberLeave (forced): %s %s", - member.Member.Name, member.Member.Addr) + s.logger.Info(fmt.Sprintf("serf: EventMemberLeave (forced): %s %s", + member.Member.Name, member.Member.Addr)) if s.config.EventCh != nil { s.config.EventCh <- MemberEvent{ Type: EventMemberLeave, @@ -1198,7 +1201,7 @@ func (s *Serf) handlePrune(member *memberState) { time.Sleep(s.config.BroadcastTimeout + s.config.LeavePropagateDelay) } - s.logger.Printf("[INFO] serf: EventMemberReap (forced): %s %s", member.Name, member.Member.Addr) + s.logger.Info(fmt.Sprintf("serf: EventMemberReap (forced): %s %s", member.Name, member.Member.Addr)) //If we are leaving or left we may be in that list of members if member.Status == StatusLeaving || member.Status == StatusLeft { @@ -1257,11 +1260,11 @@ func (s *Serf) handleUserEvent(eventMsg *messageUserEvent) bool { curTime := s.eventClock.Time() if curTime > LamportTime(len(s.eventBuffer)) && eventMsg.LTime < curTime-LamportTime(len(s.eventBuffer)) { - s.logger.Printf( - "[WARN] serf: received old event %s from time %d (current: %d)", - eventMsg.Name, - eventMsg.LTime, - s.eventClock.Time()) + s.logger.Warn( + fmt.Sprintf("serf: received old event %s from time %d (current: %d)", + eventMsg.Name, + eventMsg.LTime, + s.eventClock.Time())) return false } @@ -1316,11 +1319,11 @@ func (s *Serf) handleQuery(query *messageQuery) bool { curTime := s.queryClock.Time() if curTime > LamportTime(len(s.queryBuffer)) && query.LTime < curTime-LamportTime(len(s.queryBuffer)) { - s.logger.Printf( - "[WARN] serf: received old query %s from time %d (current: %d)", - query.Name, - query.LTime, - s.queryClock.Time()) + s.logger.Warn( + fmt.Sprintf("serf: received old query %s from time %d (current: %d)", + query.Name, + query.LTime, + s.queryClock.Time())) return false } @@ -1369,7 +1372,7 @@ func (s *Serf) handleQuery(query *messageQuery) bool { } raw, err := encodeMessage(messageQueryResponseType, &ack) if err != nil { - s.logger.Printf("[ERR] serf: failed to format ack: %v", err) + s.logger.Error(fmt.Sprintf("serf: failed to format ack: %v", err)) } else { udpAddr := net.UDPAddr{IP: query.Addr, Port: int(query.Port)} addr := memberlist.Address{ @@ -1377,10 +1380,10 @@ func (s *Serf) handleQuery(query *messageQuery) bool { Name: query.SourceNode, } if err := s.memberlist.SendToAddress(addr, raw); err != nil { - s.logger.Printf("[ERR] serf: failed to send ack: %v", err) + s.logger.Error(fmt.Sprintf("serf: failed to send ack: %v", err)) } if err := s.relayResponse(query.RelayFactor, udpAddr, query.SourceNode, &ack); err != nil { - s.logger.Printf("[ERR] serf: failed to relay ack: %v", err) + s.logger.Error(fmt.Sprintf("serf: failed to relay ack: %v", err)) } } } @@ -1410,15 +1413,15 @@ func (s *Serf) handleQueryResponse(resp *messageQueryResponse) { query, ok := s.queryResponse[resp.LTime] s.queryLock.RUnlock() if !ok { - s.logger.Printf("[WARN] serf: reply for non-running query (LTime: %d, ID: %d) From: %s", - resp.LTime, resp.ID, resp.From) + s.logger.Warn(fmt.Sprintf("serf: reply for non-running query (LTime: %d, ID: %d) From: %s", + resp.LTime, resp.ID, resp.From)) return } // Verify the ID matches if query.id != resp.ID { - s.logger.Printf("[WARN] serf: query reply ID mismatch (Local: %d, Response: %d)", - query.id, resp.ID) + s.logger.Warn(fmt.Sprintf("serf: query reply ID mismatch (Local: %d, Response: %d)", + query.id, resp.ID)) return } @@ -1438,7 +1441,7 @@ func (s *Serf) handleQueryResponse(resp *messageQueryResponse) { metrics.IncrCounterWithLabels([]string{"serf", "query_acks"}, 1, s.metricLabels) err := query.sendAck(resp) if err != nil { - s.logger.Printf("[WARN] %v", err) + s.logger.Warn(err.Error()) } } else { // Exit early if this is a duplicate response @@ -1450,7 +1453,7 @@ func (s *Serf) handleQueryResponse(resp *messageQueryResponse) { metrics.IncrCounterWithLabels([]string{"serf", "query_responses"}, 1, s.metricLabels) err := query.sendResponse(NodeResponse{From: resp.From, Payload: resp.Payload}) if err != nil { - s.logger.Printf("[WARN] %v", err) + s.logger.Warn(err.Error()) } } } @@ -1461,14 +1464,14 @@ func (s *Serf) handleQueryResponse(resp *messageQueryResponse) { func (s *Serf) handleNodeConflict(existing, other *memberlist.Node) { // Log a basic warning if the node is not us... if existing.Name != s.config.NodeName { - s.logger.Printf("[WARN] serf: Name conflict for '%s' both %s:%d and %s:%d are claiming", - existing.Name, existing.Addr, existing.Port, other.Addr, other.Port) + s.logger.Warn(fmt.Sprintf("serf: Name conflict for '%s' both %s:%d and %s:%d are claiming", + existing.Name, existing.Addr, existing.Port, other.Addr, other.Port)) return } // The current node is conflicting! This is an error - s.logger.Printf("[ERR] serf: Node name conflicts with another node at %s:%d. Names must be unique! (Resolution enabled: %v)", - other.Addr, other.Port, s.config.EnableNameConflictResolution) + s.logger.Error(fmt.Sprintf("serf: Node name conflicts with another node at %s:%d. Names must be unique! (Resolution enabled: %v)", + other.Addr, other.Port, s.config.EnableNameConflictResolution)) // If automatic resolution is enabled, kick off the resolution if s.config.EnableNameConflictResolution { @@ -1487,7 +1490,7 @@ func (s *Serf) resolveNodeConflict() { payload := []byte(s.config.NodeName) resp, err := s.Query(qName, payload, nil) if err != nil { - s.logger.Printf("[ERR] serf: Failed to start name resolution query: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to start name resolution query: %v", err)) return } @@ -1499,12 +1502,12 @@ func (s *Serf) resolveNodeConflict() { for r := range respCh { // Decode the response if len(r.Payload) < 1 || messageType(r.Payload[0]) != messageConflictResponseType { - s.logger.Printf("[ERR] serf: Invalid conflict query response type: %v", r.Payload) + s.logger.Error(fmt.Sprintf("serf: Invalid conflict query response type: %v", r.Payload)) continue } var member Member if err := decodeMessage(r.Payload[1:], &member); err != nil { - s.logger.Printf("[ERR] serf: Failed to decode conflict query response: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to decode conflict query response: %v", err)) continue } @@ -1518,20 +1521,20 @@ func (s *Serf) resolveNodeConflict() { // Query over, determine if we should live majority := (responses / 2) + 1 if matching >= majority { - s.logger.Printf("[INFO] serf: majority in name conflict resolution [%d / %d]", - matching, responses) + s.logger.Info(fmt.Sprintf("serf: majority in name conflict resolution [%d / %d]", + matching, responses)) return } // Since we lost the vote, we need to exit - s.logger.Printf("[WARN] serf: minority in name conflict resolution, quiting [%d / %d]", + s.logger.Warn("serf: minority in name conflict resolution, quiting [%d / %d]", matching, responses) if err := s.Shutdown(); err != nil { - s.logger.Printf("[ERR] serf: Failed to shutdown: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to shutdown: %v", err)) } } -//eraseNode takes a node completely out of the member list +// eraseNode takes a node completely out of the member list func (s *Serf) eraseNode(m *memberState) { // Delete from members delete(s.members, m.Name) @@ -1611,7 +1614,7 @@ func (s *Serf) reap(old []*memberState, now time.Time, timeout time.Duration) [] i-- // Delete from members and send out event - s.logger.Printf("[INFO] serf: EventMemberReap: %s", m.Name) + s.logger.Info(fmt.Sprintf("serf: EventMemberReap: %s", m.Name)) s.eraseNode(m) } @@ -1643,7 +1646,7 @@ func (s *Serf) reconnect() { prob := numFailed / numAlive if rand.Float32() > prob { s.memberLock.RUnlock() - s.logger.Printf("[DEBUG] serf: forgoing reconnect for random throttling") + s.logger.Debug("serf: forgoing reconnect for random throttling") return } @@ -1653,7 +1656,7 @@ func (s *Serf) reconnect() { // Format the addr addr := net.UDPAddr{IP: mem.Addr, Port: int(mem.Port)} - s.logger.Printf("[INFO] serf: attempting reconnect to %v %s", mem.Name, addr.String()) + s.logger.Info("serf: attempting reconnect to %v %s", mem.Name, addr.String()) joinAddr := addr.String() if mem.Name != "" { @@ -1690,11 +1693,11 @@ func (s *Serf) checkQueueDepth(name string, queue *memberlist.TransmitLimitedQue numq := queue.NumQueued() metrics.AddSampleWithLabels([]string{"serf", "queue", name}, float32(numq), s.metricLabels) if numq >= s.config.QueueDepthWarning { - s.logger.Printf("[WARN] serf: %s queue depth: %d", name, numq) + s.logger.Warn(fmt.Sprintf("serf: %s queue depth: %d", name, numq)) } if max := s.getQueueMax(); numq > max { - s.logger.Printf("[WARN] serf: %s queue depth (%d) exceeds limit (%d), dropping messages!", - name, numq, max) + s.logger.Warn(fmt.Sprintf("serf: %s queue depth (%d) exceeds limit (%d), dropping messages!", + name, numq, max)) queue.Prune(max) } case <-s.shutdownCh: @@ -1772,14 +1775,14 @@ func (s *Serf) handleRejoin(previous []*PreviousNode) { joinAddr = prev.Name + "/" + prev.Addr } - s.logger.Printf("[INFO] serf: Attempting re-join to previously known node: %s", prev) + s.logger.Info(fmt.Sprintf("serf: Attempting re-join to previously known node: %s", prev)) _, err := s.memberlist.Join([]string{joinAddr}) if err == nil { - s.logger.Printf("[INFO] serf: Re-joined to previously known node: %s", prev) + s.logger.Info(fmt.Sprintf("serf: Re-joined to previously known node: %s", prev)) return } } - s.logger.Printf("[WARN] serf: Failed to re-join any previously known node") + s.logger.Warn("serf: Failed to re-join any previously known node") } // encodeTags is used to encode a tag map @@ -1814,7 +1817,7 @@ func (s *Serf) decodeTags(buf []byte) map[string]string { r := bytes.NewReader(buf[1:]) dec := codec.NewDecoder(r, &codec.MsgpackHandle{}) if err := dec.Decode(&tags); err != nil { - s.logger.Printf("[ERR] serf: Failed to decode tags: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to decode tags: %v", err)) } return tags } diff --git a/serf/serf_test.go b/serf/serf_test.go index 0dada60fb..9662ce88e 100644 --- a/serf/serf_test.go +++ b/serf/serf_test.go @@ -9,7 +9,6 @@ import ( "encoding/base64" "fmt" "io/ioutil" - "log" "net" "os" "path/filepath" @@ -21,6 +20,7 @@ import ( "testing" "time" + "github.com/hashicorp/go-hclog" "github.com/hashicorp/go-msgpack/codec" "github.com/hashicorp/memberlist" "github.com/hashicorp/serf/coordinate" @@ -59,8 +59,10 @@ func testConfig(t *testing.T, ip net.IP) *Config { config.TombstoneTimeout = 1 * time.Microsecond if t != nil { - config.Logger = log.New(os.Stderr, "test["+t.Name()+"]: ", log.LstdFlags) - config.MemberlistConfig.Logger = config.Logger + config.Logger = hclog.Default().Named("test[" + t.Name() + "]: ") + config.MemberlistConfig.Logger = config.Logger.StandardLogger(&hclog.StandardLoggerOptions{ + InferLevels: true, + }) } return config diff --git a/serf/snapshot.go b/serf/snapshot.go index 28bb668e6..b03a15665 100644 --- a/serf/snapshot.go +++ b/serf/snapshot.go @@ -6,7 +6,6 @@ package serf import ( "bufio" "fmt" - "log" "math/rand" "net" "os" @@ -15,6 +14,7 @@ import ( "time" "github.com/armon/go-metrics" + "github.com/hashicorp/go-hclog" ) /* @@ -72,7 +72,7 @@ type Snapshotter struct { lastQueryClock LamportTime leaveCh chan struct{} leaving bool - logger *log.Logger + logger hclog.Logger minCompactSize int64 path string offset int64 @@ -103,7 +103,7 @@ func (p PreviousNode) String() string { func NewSnapshotter(path string, minCompactSize int, rejoinAfterLeave bool, - logger *log.Logger, + logger hclog.Logger, clock *LamportClock, outCh chan<- Event, shutdownCh <-chan struct{}) (chan<- Event, *Snapshotter, error) { @@ -262,7 +262,7 @@ func (s *Snapshotter) stream() { case *Query: s.processQuery(typed) default: - s.logger.Printf("[ERR] serf: Unknown event to snapshot: %#v", e) + s.logger.Error(fmt.Sprintf("serf: Unknown event to snapshot: %#v", e)) } } @@ -277,10 +277,10 @@ func (s *Snapshotter) stream() { } s.tryAppend("leave\n") if err := s.buffered.Flush(); err != nil { - s.logger.Printf("[ERR] serf: failed to flush leave to snapshot: %v", err) + s.logger.Error(fmt.Sprintf("serf: failed to flush leave to snapshot: %v", err)) } if err := s.fh.Sync(); err != nil { - s.logger.Printf("[ERR] serf: failed to sync leave to snapshot: %v", err) + s.logger.Error(fmt.Sprintf("serf: failed to sync leave to snapshot: %v", err)) } case e := <-s.streamCh: @@ -310,10 +310,10 @@ func (s *Snapshotter) stream() { } if err := s.buffered.Flush(); err != nil { - s.logger.Printf("[ERR] serf: failed to flush snapshot: %v", err) + s.logger.Error(fmt.Sprintf("serf: failed to flush snapshot: %v", err)) } if err := s.fh.Sync(); err != nil { - s.logger.Printf("[ERR] serf: failed to sync snapshot: %v", err) + s.logger.Error(fmt.Sprintf("serf: failed to sync snapshot: %v", err)) } s.fh.Close() close(s.waitCh) @@ -377,16 +377,16 @@ func (s *Snapshotter) processQuery(q *Query) { // tryAppend will invoke append line but will not return an error func (s *Snapshotter) tryAppend(l string) { if err := s.appendLine(l); err != nil { - s.logger.Printf("[ERR] serf: Failed to update snapshot: %v", err) + s.logger.Error(fmt.Sprintf("serf: Failed to update snapshot: %v", err)) now := time.Now() if now.Sub(s.lastAttemptedCompaction) > snapshotErrorRecoveryInterval { s.lastAttemptedCompaction = now - s.logger.Printf("[INFO] serf: Attempting compaction to recover from error...") + s.logger.Info("serf: Attempting compaction to recover from error...") err = s.compact() if err != nil { - s.logger.Printf("[ERR] serf: Compaction failed, will reattempt after %v: %v", snapshotErrorRecoveryInterval, err) + s.logger.Error(fmt.Sprintf("serf: Compaction failed, will reattempt after %v: %v", snapshotErrorRecoveryInterval, err)) } else { - s.logger.Printf("[INFO] serf: Finished compaction, successfully recovered from error state") + s.logger.Info("serf: Finished compaction, successfully recovered from error state") } } } @@ -564,7 +564,7 @@ func (s *Snapshotter) replay() error { info := strings.TrimPrefix(line, "alive: ") addrIdx := strings.LastIndex(info, " ") if addrIdx == -1 { - s.logger.Printf("[WARN] serf: Failed to parse address: %v", line) + s.logger.Warn(fmt.Sprintf("serf: Failed to parse address: %v", line)) continue } addr := info[addrIdx+1:] @@ -579,7 +579,7 @@ func (s *Snapshotter) replay() error { timeStr := strings.TrimPrefix(line, "clock: ") timeInt, err := strconv.ParseUint(timeStr, 10, 64) if err != nil { - s.logger.Printf("[WARN] serf: Failed to convert clock time: %v", err) + s.logger.Warn(fmt.Sprintf("serf: Failed to convert clock time: %v", err)) continue } s.lastClock = LamportTime(timeInt) @@ -588,7 +588,7 @@ func (s *Snapshotter) replay() error { timeStr := strings.TrimPrefix(line, "event-clock: ") timeInt, err := strconv.ParseUint(timeStr, 10, 64) if err != nil { - s.logger.Printf("[WARN] serf: Failed to convert event clock time: %v", err) + s.logger.Warn(fmt.Sprintf("serf: Failed to convert event clock time: %v", err)) continue } s.lastEventClock = LamportTime(timeInt) @@ -597,7 +597,7 @@ func (s *Snapshotter) replay() error { timeStr := strings.TrimPrefix(line, "query-clock: ") timeInt, err := strconv.ParseUint(timeStr, 10, 64) if err != nil { - s.logger.Printf("[WARN] serf: Failed to convert query clock time: %v", err) + s.logger.Warn(fmt.Sprintf("serf: Failed to convert query clock time: %v", err)) continue } s.lastQueryClock = LamportTime(timeInt) @@ -607,7 +607,7 @@ func (s *Snapshotter) replay() error { } else if line == "leave" { // Ignore a leave if we plan on re-joining if s.rejoinAfterLeave { - s.logger.Printf("[INFO] serf: Ignoring previous leave in snapshot") + s.logger.Info("serf: Ignoring previous leave in snapshot") continue } s.aliveNodes = make(map[string]string) @@ -619,7 +619,7 @@ func (s *Snapshotter) replay() error { // Skip comment lines } else { - s.logger.Printf("[WARN] serf: Unrecognized snapshot line: %v", line) + s.logger.Warn(fmt.Sprintf("serf: Unrecognized snapshot line: %v", line)) } } diff --git a/serf/snapshot_test.go b/serf/snapshot_test.go index bab5bc76e..854a3f6ef 100644 --- a/serf/snapshot_test.go +++ b/serf/snapshot_test.go @@ -6,11 +6,12 @@ package serf import ( "fmt" "io/ioutil" - "log" "os" "reflect" "testing" "time" + + "github.com/hashicorp/go-hclog" ) func TestSnapshotter(t *testing.T) { @@ -23,7 +24,7 @@ func TestSnapshotter(t *testing.T) { clock := new(LamportClock) outCh := make(chan Event, 64) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + logger := hclog.Default() inCh, snap, err := NewSnapshotter(td+"snap", snapshotSizeLimit, false, logger, clock, outCh, stopCh) if err != nil { @@ -176,7 +177,7 @@ func TestSnapshotter_forceCompact(t *testing.T) { clock := new(LamportClock) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + logger := hclog.Default() // Create a very low limit inCh, snap, err := NewSnapshotter(td+"snap", 1024, false, @@ -240,7 +241,7 @@ func TestSnapshotter_leave(t *testing.T) { clock := new(LamportClock) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + logger := hclog.Default() inCh, snap, err := NewSnapshotter(td+"snap", snapshotSizeLimit, false, logger, clock, nil, stopCh) if err != nil { @@ -321,7 +322,7 @@ func TestSnapshotter_leave_rejoin(t *testing.T) { clock := new(LamportClock) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + logger := hclog.Default() inCh, snap, err := NewSnapshotter(td+"snap", snapshotSizeLimit, true, logger, clock, nil, stopCh) if err != nil { @@ -404,7 +405,7 @@ func TestSnapshotter_slowDiskNotBlockingEventCh(t *testing.T) { clock := new(LamportClock) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + logger := hclog.Default() outCh := make(chan Event, 1024) inCh, snap, err := NewSnapshotter(td+"snap", snapshotSizeLimit, true, @@ -490,7 +491,7 @@ func TestSnapshotter_blockedUpstreamNotBlockingMemberlist(t *testing.T) { clock := new(LamportClock) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + logger := hclog.Default() // OutCh is unbuffered simulating a slow upstream outCh := make(chan Event) diff --git a/testutil/retry/retry.go b/testutil/retry/retry.go index f75940cf3..2914b6ade 100644 --- a/testutil/retry/retry.go +++ b/testutil/retry/retry.go @@ -5,14 +5,13 @@ // // A sample retry operation looks like this: // -// func TestX(t *testing.T) { -// retry.Run(t, func(r *retry.R) { -// if err := foo(); err != nil { -// r.Fatal("f: ", err) -// } -// }) -// } -// +// func TestX(t *testing.T) { +// retry.Run(t, func(r *retry.R) { +// if err := foo(); err != nil { +// r.Fatal("f: ", err) +// } +// }) +// } package retry import ( diff --git a/testutil/testlog.go b/testutil/testlog.go index f0e8f10ec..d88ce0d9e 100644 --- a/testutil/testlog.go +++ b/testutil/testlog.go @@ -5,17 +5,24 @@ package testutil import ( "io" - "log" "strings" "testing" + + "github.com/hashicorp/go-hclog" ) -func TestLogger(t testing.TB) *log.Logger { - return log.New(&testWriter{t}, "test: ", log.LstdFlags) +func TestLogger(t testing.TB) hclog.Logger { + return hclog.New(&hclog.LoggerOptions{ + Output: &testWriter{t}, + Name: "test: ", + }) } -func TestLoggerWithName(t testing.TB, name string) *log.Logger { - return log.New(&testWriter{t}, "test["+name+"]: ", log.LstdFlags) +func TestLoggerWithName(t testing.TB, name string) hclog.Logger { + return hclog.New(&hclog.LoggerOptions{ + Output: &testWriter{t}, + Name: "test[" + name + "]: ", + }) } func TestWriter(t testing.TB) io.Writer {