From cfc4387baa7b5a396bfd07a52d830c77f5a38a7f Mon Sep 17 00:00:00 2001 From: Luke Thorne Date: Sun, 20 Aug 2023 17:45:09 -0500 Subject: [PATCH 1/3] convert to slog --- .gitignore | 1 + client/rpc_client.go | 14 +- cmd/serf/command/agent/agent.go | 105 +++++------ cmd/serf/command/agent/agent_test.go | 8 +- cmd/serf/command/agent/command.go | 31 +++- cmd/serf/command/agent/command_test.go | 2 +- cmd/serf/command/agent/config.go | 22 +-- cmd/serf/command/agent/event_handler_test.go | 20 +-- cmd/serf/command/agent/invoke.go | 2 +- cmd/serf/command/agent/ipc.go | 64 ++++--- cmd/serf/command/agent/ipc_event_stream.go | 17 +- .../command/agent/ipc_event_stream_test.go | 17 +- cmd/serf/command/agent/ipc_log_stream.go | 15 +- cmd/serf/command/agent/ipc_log_stream_test.go | 9 +- .../agent/ipc_query_response_stream.go | 15 +- cmd/serf/command/agent/log_levels.go | 13 +- cmd/serf/command/agent/log_writer.go | 2 +- cmd/serf/command/agent/mdns.go | 24 ++- cmd/serf/command/agent/rpc_client_test.go | 12 +- cmd/serf/command/output.go | 2 +- coordinate/config.go | 13 +- go.mod | 9 +- go.sum | 30 ++-- serf/coalesce_member_test.go | 16 +- serf/config.go | 11 +- serf/delegate.go | 40 +++-- serf/delegate_test.go | 4 +- serf/internal_query.go | 63 +++---- serf/internal_query_test.go | 33 +++- serf/keymanager.go | 4 +- serf/merge_delegate.go | 4 +- serf/ping_delegate.go | 14 +- serf/query.go | 10 +- serf/serf.go | 165 +++++++++--------- serf/serf_test.go | 29 +-- serf/snapshot.go | 37 ++-- serf/snapshot_test.go | 74 ++++---- testutil/retry/retry.go | 15 +- 38 files changed, 546 insertions(+), 420 deletions(-) diff --git a/.gitignore b/.gitignore index 83e4559a4..9911dddc3 100644 --- a/.gitignore +++ b/.gitignore @@ -1,6 +1,7 @@ # Platform .DS_Store /.idea +/.vscode *.iml # Compiled Object files, Static and Dynamic libs (Shared Objects) diff --git a/client/rpc_client.go b/client/rpc_client.go index 6173cc00e..5c4dd467e 100644 --- a/client/rpc_client.go +++ b/client/rpc_client.go @@ -145,9 +145,9 @@ func ClientFromConfig(c *Config) (*RPCClient, error) { shutdownCh: make(chan struct{}), } client.dec = codec.NewDecoder(client.reader, - &codec.MsgpackHandle{RawToString: true, WriteExt: true}) + &codec.MsgpackHandle{WriteExt: true}) client.enc = codec.NewEncoder(client.writer, - &codec.MsgpackHandle{RawToString: true, WriteExt: true}) + &codec.MsgpackHandle{WriteExt: true}) go client.listen() // Do the initial handshake @@ -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, @@ -440,7 +440,7 @@ func (mh *monitorHandler) Cleanup() { if !mh.closed { if !mh.init { mh.init = true - mh.initCh <- fmt.Errorf("Stream closed") + mh.initCh <- fmt.Errorf("stream closed") } if mh.logCh != nil { close(mh.logCh) @@ -522,7 +522,7 @@ func (sh *streamHandler) Cleanup() { if !sh.closed { if !sh.init { sh.init = true - sh.initCh <- fmt.Errorf("Stream closed") + sh.initCh <- fmt.Errorf("stream closed") } if sh.eventCh != nil { close(sh.eventCh) @@ -623,7 +623,7 @@ func (qh *queryHandler) Cleanup() { if !qh.closed { if !qh.init { qh.init = true - qh.initCh <- fmt.Errorf("Stream closed") + qh.initCh <- fmt.Errorf("stream closed") } if qh.ackCh != nil { close(qh.ackCh) diff --git a/cmd/serf/command/agent/agent.go b/cmd/serf/command/agent/agent.go index e9c319b0f..f32b4e382 100644 --- a/cmd/serf/command/agent/agent.go +++ b/cmd/serf/command/agent/agent.go @@ -4,12 +4,12 @@ package agent import ( + "context" "encoding/base64" "encoding/json" "fmt" "io" - "io/ioutil" - "log" + "log/slog" "os" "strings" "sync" @@ -37,7 +37,7 @@ type Agent struct { eventHandlersLock sync.Mutex // logger instance wraps the logOutput - logger *log.Logger + logger *slog.Logger // This is the underlying Serf we are wrapping serf *serf.Serf @@ -59,6 +59,15 @@ func Create(agentConf *Config, conf *serf.Config, logOutput io.Writer) (*Agent, conf.MemberlistConfig.LogOutput = logOutput conf.MemberlistConfig.EnableCompression = agentConf.EnableCompression conf.LogOutput = logOutput + var logLevel *slog.Level + if err := logLevel.UnmarshalText([]byte(agentConf.LogLevel)); err != nil { + return nil, fmt.Errorf("error parsing log level: %s", err) + } + handlerOpts := &slog.HandlerOptions{ + AddSource: logLevel.Level() <= slog.LevelDebug, + Level: logLevel, + } + handler := slog.NewTextHandler(logOutput, handlerOpts) // Create a channel to listen for events from Serf eventCh := make(chan serf.Event, 64) @@ -70,7 +79,7 @@ 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), + logger: slog.New(handler).WithGroup("agent"), shutdownCh: make(chan struct{}), } @@ -95,12 +104,12 @@ 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.LogAttrs(context.TODO(), slog.LevelInfo, "Serf agent starting") // Create serf first serf, err := serf.Create(a.conf) if err != nil { - return fmt.Errorf("Error creating Serf: %s", err) + return fmt.Errorf("error creating Serf: %s", err) } a.serf = serf @@ -115,7 +124,7 @@ func (a *Agent) Leave() error { return nil } - a.logger.Println("[INFO] agent: requesting graceful leave from Serf") + a.logger.LogAttrs(context.TODO(), slog.LevelInfo, "requesting graceful leave from Serf") return a.serf.Leave() } @@ -133,13 +142,13 @@ func (a *Agent) Shutdown() error { goto EXIT } - a.logger.Println("[INFO] agent: requesting serf shutdown") + a.logger.LogAttrs(context.TODO(), slog.LevelInfo, "requesting serf shutdown") if err := a.serf.Shutdown(); err != nil { return err } EXIT: - a.logger.Println("[INFO] agent: shutdown complete") + a.logger.LogAttrs(context.TODO(), slog.LevelInfo, "shutdown complete") a.shutdown = true close(a.shutdownCh) return nil @@ -163,24 +172,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.LogAttrs(context.TODO(), slog.LevelInfo, "joining", slog.String("addresses", strings.Join(addrs, ",")), slog.Bool("replay", replay)) ignoreOld := !replay n, err = a.serf.Join(addrs, ignoreOld) if n > 0 { - a.logger.Printf("[INFO] agent: joined: %d nodes", n) + a.logger.LogAttrs(context.TODO(), slog.LevelInfo, "joined", slog.Int("nodes", n)) } if err != nil { - a.logger.Printf("[WARN] agent: error joining: %v", err) + a.logger.LogAttrs(context.TODO(), slog.LevelWarn, "error joining", slog.String("error", err.Error())) } 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.LogAttrs(context.TODO(), slog.LevelInfo, "Force leaving", slog.String("node", node)) err := a.serf.RemoveFailedNode(node) if err != nil { - a.logger.Printf("[WARN] agent: failed to remove node: %v", err) + a.logger.LogAttrs(context.TODO(), slog.LevelWarn, "failed to remove node", slog.String("error", err.Error())) } return err } @@ -188,21 +197,20 @@ 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.LogAttrs(context.TODO(), slog.LevelInfo, "Force leaving (prune)", slog.String("node", node)) err := a.serf.RemoveFailedNodePrune(node) if err != nil { - a.logger.Printf("[WARN] agent: failed to remove node (prune): %v", err) + a.logger.LogAttrs(context.TODO(), slog.LevelWarn, "failed to remove (prune)", slog.String("error", err.Error())) } 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.LogAttrs(context.TODO(), slog.LevelDebug, "Requesting user event send", slog.String("name", name), slog.Bool("coalesce", coalesce), slog.String("payload", 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.LogAttrs(context.TODO(), slog.LevelWarn, "failed to send user event", slog.String("error", err.Error())) } return err } @@ -213,14 +221,13 @@ func (a *Agent) Query(name string, payload []byte, params *serf.QueryParam) (*se if strings.HasPrefix(name, serf.InternalQueryPrefix) { // Allow the special "ping" query if name != serf.InternalQueryPrefix+"ping" || payload != nil { - return nil, fmt.Errorf("Queries cannot contain the '%s' prefix", serf.InternalQueryPrefix) + 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.LogAttrs(context.TODO(), slog.LevelDebug, "Requesting query send", slog.String("name", name), slog.String("payload", 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.LogAttrs(context.TODO(), slog.LevelWarn, "failed to start user query", slog.String("error", err.Error())) } return resp, err } @@ -255,7 +262,7 @@ func (a *Agent) eventLoop() { for { select { case e := <-a.eventCh: - a.logger.Printf("[INFO] agent: Received event: %s", e.String()) + a.logger.LogAttrs(context.TODO(), slog.LevelDebug, "Received event", slog.String("event", e.String())) a.eventHandlersLock.Lock() handlers := a.eventHandlerList a.eventHandlersLock.Unlock() @@ -264,7 +271,7 @@ func (a *Agent) eventLoop() { } case <-serfShutdownCh: - a.logger.Printf("[WARN] agent: Serf shutdown detected, quitting") + a.logger.LogAttrs(context.TODO(), slog.LevelWarn, "Serf shutdown detected, quitting") a.Shutdown() return @@ -276,28 +283,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.LogAttrs(context.TODO(), slog.LevelInfo, "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.LogAttrs(context.TODO(), slog.LevelInfo, "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.LogAttrs(context.TODO(), slog.LevelInfo, "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.LogAttrs(context.TODO(), slog.LevelInfo, "Initiating key listing") manager := a.serf.KeyManager() return manager.ListKeys() } @@ -308,7 +315,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to update tags file", slog.String("error", err.Error())) return err } } @@ -322,19 +329,18 @@ func (a *Agent) SetTags(tags map[string]string) error { func (a *Agent) loadTagsFile(tagsFile string) error { // Avoid passing tags and using a tags file at the same time if len(a.agentConf.Tags) > 0 { - return fmt.Errorf("Tags config not allowed while using tag files") + return fmt.Errorf("tags config not allowed while using tag files") } if _, err := os.Stat(tagsFile); err == nil { - tagData, err := ioutil.ReadFile(tagsFile) + tagData, err := os.ReadFile(tagsFile) if err != nil { - return fmt.Errorf("Failed to read tags file: %s", err) + return fmt.Errorf("failed to read tags file: %s", err) } if err := json.Unmarshal(tagData, &a.conf.Tags); err != nil { - return fmt.Errorf("Failed to decode tags file: %s", err) + return fmt.Errorf("failed to decode tags file: %s", err) } - a.logger.Printf("[INFO] agent: Restored %d tag(s) from %s", - len(a.conf.Tags), tagsFile) + a.logger.LogAttrs(context.TODO(), slog.LevelInfo, "Restored tags from file", slog.Int("count", len(a.conf.Tags)), slog.String("file", tagsFile)) } // Success! @@ -345,12 +351,12 @@ func (a *Agent) loadTagsFile(tagsFile string) error { func (a *Agent) writeTagsFile(tags map[string]string) error { encoded, err := json.MarshalIndent(tags, "", " ") if err != nil { - return fmt.Errorf("Failed to encode tags: %s", err) + return fmt.Errorf("failed to encode tags: %s", err) } // Use 0600 for permissions, in case tag data is sensitive - if err = ioutil.WriteFile(a.agentConf.TagsFile, encoded, 0600); err != nil { - return fmt.Errorf("Failed to write tags file: %s", err) + if err = os.WriteFile(a.agentConf.TagsFile, encoded, 0600); err != nil { + return fmt.Errorf("failed to write tags file: %s", err) } // Success! @@ -374,7 +380,7 @@ func UnmarshalTags(tags []string) (map[string]string, error) { for _, tag := range tags { parts := strings.SplitN(tag, "=", 2) if len(parts) != 2 || len(parts[0]) == 0 { - return nil, fmt.Errorf("Invalid tag: '%s'", tag) + return nil, fmt.Errorf("invalid tag: '%s'", tag) } result[parts[0]] = parts[1] } @@ -385,7 +391,7 @@ func UnmarshalTags(tags []string) (map[string]string, error) { func (a *Agent) loadKeyringFile(keyringFile string) error { // Avoid passing an encryption key and a keyring file at the same time if len(a.agentConf.EncryptKey) > 0 { - return fmt.Errorf("Encryption key not allowed while using a keyring") + return fmt.Errorf("encryption key not allowed while using a keyring") } if _, err := os.Stat(keyringFile); err != nil { @@ -393,15 +399,15 @@ func (a *Agent) loadKeyringFile(keyringFile string) error { } // Read in the keyring file data - keyringData, err := ioutil.ReadFile(keyringFile) + keyringData, err := os.ReadFile(keyringFile) if err != nil { - return fmt.Errorf("Failed to read keyring file: %s", err) + return fmt.Errorf("failed to read keyring file: %s", err) } // Decode keyring JSON keys := make([]string, 0) if err := json.Unmarshal(keyringData, &keys); err != nil { - return fmt.Errorf("Failed to decode keyring file: %s", err) + return fmt.Errorf("failed to decode keyring file: %s", err) } // Decode base64 values @@ -409,24 +415,23 @@ func (a *Agent) loadKeyringFile(keyringFile string) error { for i, key := range keys { keyBytes, err := base64.StdEncoding.DecodeString(key) if err != nil { - return fmt.Errorf("Failed to decode key from keyring: %s", err) + return fmt.Errorf("failed to decode key from keyring: %s", err) } keysDecoded[i] = keyBytes } // Guard against empty keyring file if len(keysDecoded) == 0 { - return fmt.Errorf("Keyring file contains no keys") + return fmt.Errorf("keyring file contains no keys") } // Create the keyring keyring, err := memberlist.NewKeyring(keysDecoded, keysDecoded[0]) if err != nil { - return fmt.Errorf("Failed to restore keyring: %s", err) + 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.LogAttrs(context.TODO(), slog.LevelInfo, "Restored keyring", slog.Int("count", len(keys)), slog.String("file", keyringFile)) // Success! return nil @@ -444,7 +449,7 @@ func (a *Agent) Stats() map[string]map[string]string { } output := map[string]map[string]string{ - "agent": map[string]string{ + "agent": { "name": local.Name, }, "runtime": runtimeStats(), diff --git a/cmd/serf/command/agent/agent_test.go b/cmd/serf/command/agent/agent_test.go index cf871e06d..938cebd5d 100644 --- a/cmd/serf/command/agent/agent_test.go +++ b/cmd/serf/command/agent/agent_test.go @@ -229,10 +229,10 @@ func TestAgent_UnmarshalTags(t *testing.T) { func TestAgent_UnmarshalTagsError(t *testing.T) { tagSets := [][]string{ - []string{"="}, - []string{"=x"}, - []string{""}, - []string{"x"}, + {"="}, + {"=x"}, + {""}, + {"x"}, } for _, tagPairs := range tagSets { if _, err := UnmarshalTags(tagPairs); err == nil { diff --git a/cmd/serf/command/agent/command.go b/cmd/serf/command/agent/command.go index 3feff1e2c..7403d825c 100644 --- a/cmd/serf/command/agent/command.go +++ b/cmd/serf/command/agent/command.go @@ -4,10 +4,12 @@ package agent import ( + "context" "flag" "fmt" "io" "log" + "log/slog" "net" "os" "os/signal" @@ -45,7 +47,7 @@ type Command struct { args []string scriptHandler *ScriptEventHandler logFilter *logutils.LevelFilter - logger *log.Logger + logger *slog.Logger } var _ cli.Command = &Command{} @@ -400,8 +402,19 @@ func (c *Command) setupLoggers(config *Config) (*GatedWriter, *logWriter, io.Wri logOutput = io.MultiWriter(c.logFilter, logWriter) } + var logLevel *slog.Level + if err := logLevel.UnmarshalText([]byte(config.LogLevel)); err != nil { + c.Ui.Error(fmt.Sprintf("Error parsing log level: %s", err)) + return nil, nil, nil + } + handlerOpts := &slog.HandlerOptions{ + AddSource: logLevel.Level() <= slog.LevelDebug, + Level: logLevel, + } + handler := slog.NewTextHandler(logOutput, handlerOpts) + // Create a logger - c.logger = log.New(logOutput, "", log.LstdFlags) + c.logger = slog.New(handler) return logGate, logWriter, logOutput } @@ -424,6 +437,10 @@ func (c *Command) startAgent(config *Config, agent *Agent, // Parse the bind address information bindIP, bindPort, err := config.AddrParts(config.BindAddr) + if err != nil { + c.Ui.Error(fmt.Sprintf("Failed to parse bind address: %v", err)) + return nil + } bindAddr := &net.TCPAddr{IP: net.ParseIP(bindIP), Port: bindPort} // Start the discovery layer @@ -504,23 +521,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.LogAttrs(context.TODO(), slog.LevelInfo, "Joining cluster...", slog.Bool("ReplayOnJoin", 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.LogAttrs(context.TODO(), slog.LevelInfo, "Join completed", slog.Int("initial synced 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.LogAttrs(context.TODO(), slog.LevelError, "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.LogAttrs(context.TODO(), slog.LevelWarn, "Join failed: %v, retrying in %v", slog.String("error", err.Error()), slog.Duration("retrying in", config.RetryInterval)) time.Sleep(config.RetryInterval) } } @@ -686,7 +703,7 @@ 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 } diff --git a/cmd/serf/command/agent/command_test.go b/cmd/serf/command/agent/command_test.go index ee098b754..bbe26daf7 100644 --- a/cmd/serf/command/agent/command_test.go +++ b/cmd/serf/command/agent/command_test.go @@ -245,7 +245,7 @@ func TestCommandRun_advertiseAddr(t *testing.T) { // Check the addr and port is as advertised! m := members[0] - if bytes.Compare(m.Addr, []byte{127, 0, 0, 10}) != 0 { + if !bytes.Equal(m.Addr, []byte{127, 0, 0, 10}) { t.Fatalf("bad: %#v", m) } if m.Port != 12345 { diff --git a/cmd/serf/command/agent/config.go b/cmd/serf/command/agent/config.go index 47eb39f5a..bb48c32d3 100644 --- a/cmd/serf/command/agent/config.go +++ b/cmd/serf/command/agent/config.go @@ -31,7 +31,7 @@ func DefaultConfig() *Config { AdvertiseAddr: "", LogLevel: "INFO", RPCAddr: "127.0.0.1:7373", - Protocol: serf.ProtocolVersionMax, + Protocol: int(serf.ProtocolVersionMax), ReplayOnJoin: false, Profile: "lan", RetryInterval: 30 * time.Second, @@ -377,7 +377,7 @@ func MergeConfig(a, b *Config) *Config { if b.Role != "" { result.Role = b.Role } - if b.DisableCoordinates == true { + if b.DisableCoordinates { result.DisableCoordinates = true } if b.Tags != nil { @@ -409,7 +409,7 @@ func MergeConfig(a, b *Config) *Config { if b.RPCAuthKey != "" { result.RPCAuthKey = b.RPCAuthKey } - if b.ReplayOnJoin != false { + if b.ReplayOnJoin { result.ReplayOnJoin = b.ReplayOnJoin } if b.Profile != "" { @@ -418,10 +418,10 @@ func MergeConfig(a, b *Config) *Config { if b.SnapshotPath != "" { result.SnapshotPath = b.SnapshotPath } - if b.LeaveOnTerm == true { + if b.LeaveOnTerm { result.LeaveOnTerm = true } - if b.SkipLeaveOnInt == true { + if b.SkipLeaveOnInt { result.SkipLeaveOnInt = true } if b.Discover != "" { @@ -510,13 +510,13 @@ func ReadConfigPaths(paths []string) (*Config, error) { for _, path := range paths { f, err := os.Open(path) if err != nil { - return nil, fmt.Errorf("Error reading '%s': %s", path, err) + return nil, fmt.Errorf("error reading '%s': %s", path, err) } fi, err := f.Stat() if err != nil { f.Close() - return nil, fmt.Errorf("Error reading '%s': %s", path, err) + return nil, fmt.Errorf("error reading '%s': %s", path, err) } if !fi.IsDir() { @@ -524,7 +524,7 @@ func ReadConfigPaths(paths []string) (*Config, error) { f.Close() if err != nil { - return nil, fmt.Errorf("Error decoding '%s': %s", path, err) + return nil, fmt.Errorf("error decoding '%s': %s", path, err) } result = MergeConfig(result, config) @@ -534,7 +534,7 @@ func ReadConfigPaths(paths []string) (*Config, error) { contents, err := f.Readdir(-1) f.Close() if err != nil { - return nil, fmt.Errorf("Error reading '%s': %s", path, err) + return nil, fmt.Errorf("error reading '%s': %s", path, err) } // Sort the contents, ensures lexical order @@ -554,14 +554,14 @@ func ReadConfigPaths(paths []string) (*Config, error) { subpath := filepath.Join(path, fi.Name()) f, err := os.Open(subpath) if err != nil { - return nil, fmt.Errorf("Error reading '%s': %s", subpath, err) + return nil, fmt.Errorf("error reading '%s': %s", subpath, err) } config, err := DecodeConfig(f) f.Close() if err != nil { - return nil, fmt.Errorf("Error decoding '%s': %s", subpath, err) + return nil, fmt.Errorf("error decoding '%s': %s", subpath, err) } result = MergeConfig(result, config) diff --git a/cmd/serf/command/agent/event_handler_test.go b/cmd/serf/command/agent/event_handler_test.go index faefc04a2..abbe954d1 100644 --- a/cmd/serf/command/agent/event_handler_test.go +++ b/cmd/serf/command/agent/event_handler_test.go @@ -388,44 +388,44 @@ func TestParseEventFilter(t *testing.T) { }{ { "", - []EventFilter{EventFilter{"*", ""}}, + []EventFilter{{"*", ""}}, }, { "member-join", - []EventFilter{EventFilter{"member-join", ""}}, + []EventFilter{{"member-join", ""}}, }, { "member-reap", - []EventFilter{EventFilter{"member-reap", ""}}, + []EventFilter{{"member-reap", ""}}, }, { "foo,bar", []EventFilter{ - EventFilter{"foo", ""}, - EventFilter{"bar", ""}, + {"foo", ""}, + {"bar", ""}, }, }, { "user:deploy", - []EventFilter{EventFilter{"user", "deploy"}}, + []EventFilter{{"user", "deploy"}}, }, { "foo,user:blah,bar", []EventFilter{ - EventFilter{"foo", ""}, - EventFilter{"user", "blah"}, - EventFilter{"bar", ""}, + {"foo", ""}, + {"user", "blah"}, + {"bar", ""}, }, }, { "query:load", - []EventFilter{EventFilter{"query", "load"}}, + []EventFilter{{"query", "load"}}, }, } diff --git a/cmd/serf/command/agent/invoke.go b/cmd/serf/command/agent/invoke.go index 8a4dbea0a..a43fc9d90 100644 --- a/cmd/serf/command/agent/invoke.go +++ b/cmd/serf/command/agent/invoke.go @@ -93,7 +93,7 @@ func invokeEventScript(logger *log.Logger, script string, self serf.Member, even cmd.Env = append(cmd.Env, fmt.Sprintf("SERF_QUERY_LTIME=%d", e.LTime)) go streamPayload(logger, stdin, e.Payload) default: - return fmt.Errorf("Unknown event type: %s", event.EventType().String()) + return fmt.Errorf("unknown event type: %s", event.EventType().String()) } // Start a timer to warn about slow handlers diff --git a/cmd/serf/command/agent/ipc.go b/cmd/serf/command/agent/ipc.go index c4a532d59..e4ddd5639 100644 --- a/cmd/serf/command/agent/ipc.go +++ b/cmd/serf/command/agent/ipc.go @@ -26,9 +26,10 @@ package agent import ( "bufio" + "context" "fmt" "io" - "log" + "log/slog" "net" "os" "regexp" @@ -73,16 +74,16 @@ const ( ) const ( - unsupportedCommand = "Unsupported command" - unsupportedIPCVersion = "Unsupported IPC version" - duplicateHandshake = "Handshake already performed" - handshakeRequired = "Handshake required" - monitorExists = "Monitor already exists" - invalidFilter = "Invalid event filter" - streamExists = "Stream with given sequence exists" - invalidQueryID = "No pending queries matching ID" - authRequired = "Authentication required" - invalidAuthToken = "Invalid authentication token" + unsupportedCommand = "unsupported command" + unsupportedIPCVersion = "unsupported IPC version" + duplicateHandshake = "handshake already performed" + handshakeRequired = "handshake required" + monitorExists = "monitor already exists" + invalidFilter = "invalid event filter" + streamExists = "stream with given sequence exists" + invalidQueryID = "no pending queries matching ID" + authRequired = "authentication required" + invalidAuthToken = "invalid authentication token" ) const ( @@ -245,7 +246,7 @@ type AgentIPC struct { authKey string clients map[string]*IPCClient listener net.Listener - logger *log.Logger + logger *slog.Logger logWriter *logWriter stop uint32 stopCh chan struct{} @@ -309,7 +310,7 @@ func (c *IPCClient) RegisterQuery(q *serf.Query) uint64 { id := c.nextQueryID() // Ensure the query deadline is in the future - timeout := q.Deadline().Sub(time.Now()) + timeout := time.Until(q.Deadline()) if timeout < 0 { return id } @@ -334,12 +335,21 @@ func NewAgentIPC(agent *Agent, authKey string, listener net.Listener, if logOutput == nil { logOutput = os.Stderr } + var logLevel *slog.Level + if err := logLevel.UnmarshalText([]byte(agent.agentConf.LogLevel)); err != nil { + return nil + } + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stdout, handlerOpts) ipc := &AgentIPC{ agent: agent, authKey: authKey, clients: make(map[string]*IPCClient), listener: listener, - logger: log.New(logOutput, "", log.LstdFlags), + logger: slog.New(handler).WithGroup("agent.ipc"), logWriter: logWriter, stopCh: make(chan struct{}), } @@ -378,10 +388,10 @@ func (i *AgentIPC) listen() { if i.isStopped() { return } - i.logger.Printf("[ERR] agent.ipc: Failed to accept client: %v", err) + i.logger.LogAttrs(context.TODO(), slog.LevelError, "Failed to accept client", slog.String("error", err.Error())) continue } - i.logger.Printf("[INFO] agent.ipc: Accepted client: %v", conn.RemoteAddr()) + i.logger.LogAttrs(context.TODO(), slog.LevelInfo, "Accepted client", slog.String("client", conn.RemoteAddr().String())) metrics.IncrCounterWithLabels([]string{"agent", "ipc", "accept"}, 1, nil) // Wrap the connection in a client @@ -394,9 +404,9 @@ func (i *AgentIPC) listen() { pendingQueries: make(map[uint64]*serf.Query), } client.dec = codec.NewDecoder(client.reader, - &codec.MsgpackHandle{RawToString: true, WriteExt: true}) + &codec.MsgpackHandle{WriteExt: true}) client.enc = codec.NewEncoder(client.writer, - &codec.MsgpackHandle{RawToString: true, WriteExt: true}) + &codec.MsgpackHandle{WriteExt: true}) // Register the client i.Lock() @@ -445,7 +455,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.LogAttrs(context.TODO(), slog.LevelError, "failed to decode request header", slog.String("error", err.Error())) } } return @@ -453,7 +463,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to evaluate request", slog.String("error", err.Error())) return } } @@ -475,7 +485,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.LogAttrs(context.TODO(), slog.LevelWarn, "Client sending commands before auth") respHeader := responseHeader{Seq: seq, Error: authRequired} client.Send(&respHeader, nil) return nil @@ -702,19 +712,19 @@ func (i *AgentIPC) filterMembers(members []serf.Member, tags map[string]string, for tag, expr := range tags { re, err := regexp.Compile(fmt.Sprintf("^%s$", expr)) if err != nil { - return nil, fmt.Errorf("Failed to compile regex: %v", err) + return nil, fmt.Errorf("failed to compile regex: %v", err) } tagsRe[tag] = re } statusRe, err := regexp.Compile(fmt.Sprintf("^%s$", status)) if err != nil { - return nil, fmt.Errorf("Failed to compile regex: %v", err) + return nil, fmt.Errorf("failed to compile regex: %v", err) } nameRe, err := regexp.Compile(fmt.Sprintf("^%s$", name)) if err != nil { - return nil, fmt.Errorf("Failed to compile regex: %v", err) + return nil, fmt.Errorf("failed to compile regex: %v", err) } OUTER: @@ -931,12 +941,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.LogAttrs(context.TODO(), slog.LevelInfo, "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.LogAttrs(context.TODO(), slog.LevelError, "leave failed", slog.String("error", err.Error())) } resp := responseHeader{Seq: seq, Error: errToString(err)} @@ -945,7 +955,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.LogAttrs(context.TODO(), slog.LevelError, "shutdown failed", slog.String("error", err.Error())) } return err } diff --git a/cmd/serf/command/agent/ipc_event_stream.go b/cmd/serf/command/agent/ipc_event_stream.go index ad18c4e32..6e1deb0ef 100644 --- a/cmd/serf/command/agent/ipc_event_stream.go +++ b/cmd/serf/command/agent/ipc_event_stream.go @@ -4,8 +4,9 @@ package agent import ( + "context" "fmt" - "log" + "log/slog" "github.com/hashicorp/serf/serf" ) @@ -20,16 +21,16 @@ type eventStream struct { client streamClient eventCh chan serf.Event filters []EventFilter - logger *log.Logger + logger *slog.Logger seq uint64 } -func newEventStream(client streamClient, filters []EventFilter, seq uint64, logger *log.Logger) *eventStream { +func newEventStream(client streamClient, filters []EventFilter, seq uint64, logger *slog.Logger) *eventStream { es := &eventStream{ client: client, eventCh: make(chan serf.Event, 512), filters: filters, - logger: logger, + logger: logger.WithGroup("agent.ipc"), seq: seq, } go es.stream() @@ -50,7 +51,7 @@ HANDLE: select { case es.eventCh <- e: default: - es.logger.Printf("[WARN] agent.ipc: Dropping event to %v", es.client) + es.logger.LogAttrs(context.TODO(), slog.LevelWarn, "Dropping event", slog.Any("to", es.client)) } } @@ -69,11 +70,11 @@ func (es *eventStream) stream() { case *serf.Query: err = es.sendQuery(e) default: - err = fmt.Errorf("Unknown event type: %s", event.EventType().String()) + 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to stream event to %v: %v", + slog.Any("destination", es.client), slog.String("error", err.Error())) 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..1298ac8e4 100644 --- a/cmd/serf/command/agent/ipc_event_stream_test.go +++ b/cmd/serf/command/agent/ipc_event_stream_test.go @@ -5,7 +5,7 @@ package agent import ( "bytes" - "log" + "log/slog" "net" "os" "testing" @@ -33,7 +33,12 @@ 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)) + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stdout, handlerOpts) + es := newEventStream(sc, filters, 42, slog.New(handler)) defer es.Stop() es.HandleEvent(serf.UserEvent{ @@ -51,7 +56,7 @@ func TestIPCEventStream(t *testing.T) { es.HandleEvent(serf.MemberEvent{ Type: serf.EventMemberJoin, Members: []serf.Member{ - serf.Member{ + { Name: "TestNode", Addr: net.IP([]byte{127, 0, 0, 1}), Port: 12345, @@ -96,7 +101,7 @@ func TestIPCEventStream(t *testing.T) { if obj1.Name != "foobar" { t.Fatalf("bad event: %#v", obj1) } - if bytes.Compare(obj1.Payload, []byte("test")) != 0 { + if !bytes.Equal(obj1.Payload, []byte("test")) { t.Fatalf("bad event: %#v", obj1) } if !obj1.Coalesce { @@ -111,7 +116,7 @@ func TestIPCEventStream(t *testing.T) { if mem1.Name != "TestNode" { t.Fatalf("bad member: %#v", mem1) } - if bytes.Compare(mem1.Addr, []byte{127, 0, 0, 1}) != 0 { + if !bytes.Equal(mem1.Addr, []byte{127, 0, 0, 1}) { t.Fatalf("bad member: %#v", mem1) } if mem1.Port != 12345 { @@ -152,7 +157,7 @@ func TestIPCEventStream(t *testing.T) { if obj3.Name != "deploy" { t.Fatalf("bad query: %#v", obj3) } - if bytes.Compare(obj3.Payload, []byte("test")) != 0 { + if !bytes.Equal(obj3.Payload, []byte("test")) { t.Fatalf("bad query: %#v", obj3) } diff --git a/cmd/serf/command/agent/ipc_log_stream.go b/cmd/serf/command/agent/ipc_log_stream.go index ce32bfaee..421ff3fda 100644 --- a/cmd/serf/command/agent/ipc_log_stream.go +++ b/cmd/serf/command/agent/ipc_log_stream.go @@ -4,7 +4,8 @@ package agent import ( - "log" + "context" + "log/slog" "github.com/hashicorp/logutils" ) @@ -14,17 +15,17 @@ type logStream struct { client streamClient filter *logutils.LevelFilter logCh chan string - logger *log.Logger + logger *slog.Logger seq uint64 } func newLogStream(client streamClient, filter *logutils.LevelFilter, - seq uint64, logger *log.Logger) *logStream { + seq uint64, logger *slog.Logger) *logStream { ls := &logStream{ client: client, filter: filter, logCh: make(chan string, 512), - logger: logger, + logger: logger.WithGroup("agent.ipc"), seq: seq, } go ls.stream() @@ -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.LogAttrs(context.TODO(), slog.LevelWarn, "Dropping logs", slog.Any("client", 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.LogAttrs(context.TODO(), slog.LevelError, " Failed to stream log", + slog.Any("client", ls.client), slog.String("error", err.Error())) 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..23628739d 100644 --- a/cmd/serf/command/agent/ipc_log_stream_test.go +++ b/cmd/serf/command/agent/ipc_log_stream_test.go @@ -4,7 +4,7 @@ package agent import ( - "log" + "log/slog" "os" "testing" "time" @@ -17,7 +17,12 @@ func TestIPCLogStream(t *testing.T) { filter := LevelFilter() filter.MinLevel = logutils.LogLevel("INFO") - ls := newLogStream(sc, filter, 42, log.New(os.Stderr, "", log.LstdFlags)) + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stdout, handlerOpts) + ls := newLogStream(sc, filter, 42, slog.New(handler)) 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..a097deb85 100644 --- a/cmd/serf/command/agent/ipc_query_response_stream.go +++ b/cmd/serf/command/agent/ipc_query_response_stream.go @@ -4,7 +4,8 @@ package agent import ( - "log" + "context" + "log/slog" "time" "github.com/hashicorp/serf/serf" @@ -13,11 +14,11 @@ import ( // queryResponseStream is used to stream the query results back to a client type queryResponseStream struct { client streamClient - logger *log.Logger + logger *slog.Logger seq uint64 } -func newQueryResponseStream(client streamClient, seq uint64, logger *log.Logger) *queryResponseStream { +func newQueryResponseStream(client streamClient, seq uint64, logger *slog.Logger) *queryResponseStream { qs := &queryResponseStream{ client: client, logger: logger, @@ -29,7 +30,7 @@ func newQueryResponseStream(client streamClient, seq uint64, logger *log.Logger) // Stream is a long running routine used to stream the results of a query back to a client func (qs *queryResponseStream) Stream(resp *serf.QueryResponse) { // Setup a timer for the query ending - remaining := resp.Deadline().Sub(time.Now()) + remaining := time.Until(resp.Deadline()) done := time.After(remaining) ackCh := resp.AckCh() @@ -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.LogAttrs(context.TODO(), slog.LevelError, "Failed to stream ack to %v: %v", slog.Any("client", qs.client), slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to stream response to %v: %v", slog.Any("client", qs.client), slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to stream query end to %v: %v", slog.Any("client", qs.client), slog.String("error", err.Error())) } return } diff --git a/cmd/serf/command/agent/log_levels.go b/cmd/serf/command/agent/log_levels.go index 98468b48f..82aec2928 100644 --- a/cmd/serf/command/agent/log_levels.go +++ b/cmd/serf/command/agent/log_levels.go @@ -4,18 +4,27 @@ package agent import ( - "io/ioutil" + "io" + "log/slog" "github.com/hashicorp/logutils" ) +const ( + LevelTrace slog.Level = slog.LevelDebug - 4 + LevelDebug slog.Level = slog.LevelDebug + LevelInfo slog.Level = slog.LevelInfo + LevelWarn slog.Level = slog.LevelWarn + LevelError slog.Level = slog.LevelError +) + // LevelFilter returns a LevelFilter that is configured with the log // levels that we use. func LevelFilter() *logutils.LevelFilter { return &logutils.LevelFilter{ Levels: []logutils.LogLevel{"TRACE", "DEBUG", "INFO", "WARN", "ERR"}, MinLevel: "INFO", - Writer: ioutil.Discard, + Writer: io.Discard, } } diff --git a/cmd/serf/command/agent/log_writer.go b/cmd/serf/command/agent/log_writer.go index fc58a5710..62c769ec8 100644 --- a/cmd/serf/command/agent/log_writer.go +++ b/cmd/serf/command/agent/log_writer.go @@ -79,7 +79,7 @@ func (l *logWriter) Write(p []byte) (n int, err error) { l.logs[l.index] = string(p) l.index = (l.index + 1) % len(l.logs) - for lh, _ := range l.handlers { + for lh := range l.handlers { lh.HandleLog(string(p)) } return diff --git a/cmd/serf/command/agent/mdns.go b/cmd/serf/command/agent/mdns.go index 6ee6f5f4d..ba6df7660 100644 --- a/cmd/serf/command/agent/mdns.go +++ b/cmd/serf/command/agent/mdns.go @@ -4,10 +4,12 @@ package agent import ( + "context" "fmt" "io" - "log" + "log/slog" "net" + "os" "time" "github.com/hashicorp/mdns" @@ -23,7 +25,7 @@ const ( type AgentMDNS struct { agent *Agent discover string - logger *log.Logger + logger *slog.Logger seen map[string]struct{} server *mdns.Server replay bool @@ -58,11 +60,21 @@ func NewAgentMDNS(agent *Agent, logOutput io.Writer, replay bool, return nil, err } + var logLevel *slog.Level + if err := logLevel.UnmarshalText([]byte(agent.agentConf.LogLevel)); err != nil { + return nil, err + } + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stdout, handlerOpts) + // Initialize the AgentMDNS m := &AgentMDNS{ agent: agent, discover: discover, - logger: log.New(logOutput, "", log.LstdFlags), + logger: slog.New(handler).WithGroup("agent.mdns"), seen: make(map[string]struct{}), server: server, replay: replay, @@ -101,10 +113,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to join", slog.String("error", err.Error())) } if n > 0 { - m.logger.Printf("[INFO] agent.mdns: Joined %d hosts", n) + m.logger.LogAttrs(context.TODO(), slog.LevelInfo, "Joined hosts", slog.Int("count", n)) } // Mark all as seen @@ -128,7 +140,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to poll for new hosts", slog.String("error", err.Error())) } } diff --git a/cmd/serf/command/agent/rpc_client_test.go b/cmd/serf/command/agent/rpc_client_test.go index 29b86fe82..db25a2eae 100644 --- a/cmd/serf/command/agent/rpc_client_test.go +++ b/cmd/serf/command/agent/rpc_client_test.go @@ -104,7 +104,7 @@ WAIT: if len(m) != 2 { t.Fatalf("should have 2 members: %#v", a1.Serf().Members()) } - if findMember(t, m, a2.conf.NodeName).Status != serf.StatusFailed && time.Now().Sub(start) < 3*time.Second { + if findMember(t, m, a2.conf.NodeName).Status != serf.StatusFailed && time.Since(start) < 3*time.Second { goto WAIT } @@ -166,7 +166,7 @@ WAIT: if len(m) != 2 { t.Fatalf("should have 2 members: %#v", a1.Serf().Members()) } - if findMember(t, m, a2.conf.NodeName).Status != serf.StatusFailed && time.Now().Sub(start) < 3*time.Second { + if findMember(t, m, a2.conf.NodeName).Status != serf.StatusFailed && time.Since(start) < 3*time.Second { goto WAIT } @@ -545,7 +545,7 @@ func TestRPCClientStream_User(t *testing.T) { if e["Name"].(string) != "deploy" { t.Fatalf("bad event: %#v", e) } - if bytes.Compare(e["Payload"].([]byte), []byte("foo")) != 0 { + if !bytes.Equal(e["Payload"].([]byte), []byte("foo")) { t.Fatalf("bad event: %#v", e) } if e["Coalesce"].(bool) != false { @@ -816,7 +816,7 @@ func TestRPCClientStream_Query(t *testing.T) { if e["Name"].(string) != "deploy" { t.Fatalf("bad query: %#v", e) } - if bytes.Compare(e["Payload"].([]byte), []byte("foo")) != 0 { + if !bytes.Equal(e["Payload"].([]byte), []byte("foo")) { t.Fatalf("bad query: %#v", e) } @@ -1017,7 +1017,7 @@ func TestRPCClient_Keys(t *testing.T) { testutil.Yield() - keys, num, _, err := client.ListKeys() + keys, _, _, err := client.ListKeys() if err != nil { t.Fatalf("err: %v", err) } @@ -1055,7 +1055,7 @@ func TestRPCClient_Keys(t *testing.T) { } // New key should now appear in the list of keys - keys, num, _, err = client.ListKeys() + keys, num, _, err := client.ListKeys() if err != nil { t.Fatalf("err: %v", err) } diff --git a/cmd/serf/command/output.go b/cmd/serf/command/output.go index 46ade3247..93e5a8578 100644 --- a/cmd/serf/command/output.go +++ b/cmd/serf/command/output.go @@ -28,7 +28,7 @@ func formatOutput(data interface{}, format string) ([]byte, error) { out = data.(fmt.Stringer).String() default: - return nil, fmt.Errorf("Invalid output format \"%s\"", format) + return nil, fmt.Errorf("invalid output format \"%s\"", format) } return []byte(prepareOutput(out)), nil 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..85973a1ba 100644 --- a/go.mod +++ b/go.mod @@ -6,15 +6,14 @@ 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-msgpack v0.5.3 - github.com/hashicorp/go-multierror v1.1.0 // indirect + github.com/fatih/color v1.13.0 // indirect + github.com/hashicorp/go-msgpack v1.1.5 + github.com/hashicorp/go-multierror v1.1.1 // indirect github.com/hashicorp/go-syslog v1.0.0 - github.com/hashicorp/go-uuid v1.0.1 // indirect + github.com/hashicorp/go-uuid v1.0.2 // indirect 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..eb4c8e941 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= @@ -53,19 +53,20 @@ github.com/hashicorp/errwrap v1.0.0/go.mod h1:YH+1FKiLXxHSkmPseP+kNlulaMuP3n2brv github.com/hashicorp/go-cleanhttp v0.5.0/go.mod h1:JpRdi6/HCYpAwUzNwuwqhbovhLtngrth3wmdIIUrZ80= 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= github.com/hashicorp/go-msgpack v0.5.3/go.mod h1:ahLV/dePpqEmjfWmKiqvPkv/twdG7iPBM1vqhUKIvfM= +github.com/hashicorp/go-msgpack v1.1.5 h1:9byZdVjKTe5mce63pRVNP1L7UAmdHOTEMGehn6KvJWs= +github.com/hashicorp/go-msgpack v1.1.5/go.mod h1:gWVc3sv/wbDmR3rQsj1CAktEZzoz1YNK9NfGLXJ69/4= github.com/hashicorp/go-multierror v1.0.0/go.mod h1:dHtQlpGsu+cZNNAkkCN/P3hoUDHhCYQXV3UM06sGGrk= -github.com/hashicorp/go-multierror v1.1.0 h1:B9UzwGQJehnUY1yNrnwREHc3fGbC2xefo8g4TbElacI= -github.com/hashicorp/go-multierror v1.1.0/go.mod h1:spPvp8C1qA32ftKqdAHm4hHTbPw+vmowP0z+KUhOZdA= +github.com/hashicorp/go-multierror v1.1.1 h1:H5DkEtf6CXdFp0N0Em5UCwQpXMWke8IA0+lD48awMYo= +github.com/hashicorp/go-multierror v1.1.1/go.mod h1:iw975J/qwKPdAO1clOe2L8331t/9/fmwbPZ6JB6eMoM= github.com/hashicorp/go-retryablehttp v0.5.3/go.mod h1:9B5zBasrRhHXnJnui7y6sL7es7NDiJgTc6Er0maI1Xs= github.com/hashicorp/go-sockaddr v1.0.0 h1:GeH6tui99pF4NJgfnhp+L6+FfobzVW3Ah46sLo0ICXs= github.com/hashicorp/go-sockaddr v1.0.0/go.mod h1:7Xibr9yA9JjQq1JpNB2Vw7kxv8xerXegt+ozgdvDeDU= github.com/hashicorp/go-syslog v1.0.0 h1:KaodqZuhUoZereWVIYmpUgZysurB1kBLX2j0MwMrUAE= github.com/hashicorp/go-syslog v1.0.0/go.mod h1:qPfqrKkXGihmCqbJM2mZgkZGvKG1dFdvsLplgctolz4= github.com/hashicorp/go-uuid v1.0.0/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro= -github.com/hashicorp/go-uuid v1.0.1 h1:fv1ep09latC32wFoVwnqcnKJGnMSdBanPczbHAYm1BE= -github.com/hashicorp/go-uuid v1.0.1/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro= +github.com/hashicorp/go-uuid v1.0.2 h1:cfejS+Tpcp13yd5nYHWDI6qVCny6wyX2Mt5SGur2IGE= +github.com/hashicorp/go-uuid v1.0.2/go.mod h1:6SBZvOh/SIDV7/2o3Jml5SYk/TvGqwFJ/bN7x4byOro= github.com/hashicorp/golang-lru v0.5.0 h1:CL2msUPvZTLb5O648aiLNJw3hnBxN2+1Jq8rCOH9wdo= github.com/hashicorp/golang-lru v0.5.0/go.mod h1:/m3WP610KZHVQ1SGc6re/UDhFvYD7pJ4Ao+sR/qLZy8= github.com/hashicorp/logutils v1.0.0 h1:dLEQVugN8vlakKOUE3ihGLTZJRB4j+M2cdTm/ORI65Y= @@ -90,14 +91,12 @@ 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 h1:sqDoxXbdeALODt0DAeJCVp38ps9ZogZEAXjus69YV3U= +github.com/mattn/go-colorable v0.1.9/go.mod h1:u6P/XSegPjTcexA+o6vUJrdnUu04hMope9wVRipJSqc= 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= @@ -162,6 +161,7 @@ golang.org/x/crypto v0.0.0-20200414173820-0848c9571904/go.mod h1:LzIPMQfyMNhhGPh golang.org/x/crypto v0.0.0-20200820211705-5c72a883971a h1:vclmkQCjlDX5OydZ9wv8rBCcS0QyQY66Mpf/7BZbInM= golang.org/x/crypto v0.0.0-20200820211705-5c72a883971a/go.mod h1:LzIPMQfyMNhhGPhUkYOs5KpL4U8rLKemX1yGLhDgUto= golang.org/x/net v0.0.0-20181114220301-adae6a3d119a/go.mod h1:mL1N/T3taQHkDXs73rZJwtUhF3w3ftmwwsq0BUmARs4= +golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= golang.org/x/net v0.0.0-20190404232315-eb5bcb51f2a3/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg= golang.org/x/net v0.0.0-20190613194153-d28f0bde5980/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s= @@ -178,18 +178,17 @@ 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-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= @@ -198,6 +197,7 @@ golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk= golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ= golang.org/x/tools v0.0.0-20180917221912-90fa682c2a6e/go.mod h1:n7NCudcB/nEzxVGmLbDWY5pfWTLqBcC2KZ6jyYvM4mQ= +golang.org/x/tools v0.0.0-20190424220101-1e8e1cfdf96b/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q= golang.org/x/tools v0.0.0-20190907020128-2ca718005c18/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo= golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= golang.org/x/xerrors v0.0.0-20191204190536-9bdfabe68543/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0= diff --git a/serf/coalesce_member_test.go b/serf/coalesce_member_test.go index 83fdbcf61..f37dabb76 100644 --- a/serf/coalesce_member_test.go +++ b/serf/coalesce_member_test.go @@ -26,27 +26,27 @@ func TestMemberEventCoalesce_Basic(t *testing.T) { send := []Event{ MemberEvent{ Type: EventMemberJoin, - Members: []Member{Member{Name: "foo"}}, + Members: []Member{{Name: "foo"}}, }, MemberEvent{ Type: EventMemberLeave, - Members: []Member{Member{Name: "foo"}}, + Members: []Member{{Name: "foo"}}, }, MemberEvent{ Type: EventMemberLeave, - Members: []Member{Member{Name: "bar"}}, + Members: []Member{{Name: "bar"}}, }, MemberEvent{ Type: EventMemberUpdate, - Members: []Member{Member{Name: "zip", Tags: map[string]string{"role": "foo"}}}, + Members: []Member{{Name: "zip", Tags: map[string]string{"role": "foo"}}}, }, MemberEvent{ Type: EventMemberUpdate, - Members: []Member{Member{Name: "zip", Tags: map[string]string{"role": "bar"}}}, + Members: []Member{{Name: "zip", Tags: map[string]string{"role": "bar"}}}, }, MemberEvent{ Type: EventMemberReap, - Members: []Member{Member{Name: "dead"}}, + Members: []Member{{Name: "dead"}}, }, } @@ -134,7 +134,7 @@ func TestMemberEventCoalesce_TagUpdate(t *testing.T) { inCh <- MemberEvent{ Type: EventMemberUpdate, - Members: []Member{Member{Name: "foo", Tags: map[string]string{"role": "foo"}}}, + Members: []Member{{Name: "foo", Tags: map[string]string{"role": "foo"}}}, } time.Sleep(30 * time.Millisecond) @@ -152,7 +152,7 @@ func TestMemberEventCoalesce_TagUpdate(t *testing.T) { // last event was an update inCh <- MemberEvent{ Type: EventMemberUpdate, - Members: []Member{Member{Name: "foo", Tags: map[string]string{"role": "bar"}}}, + Members: []Member{{Name: "foo", Tags: map[string]string{"role": "bar"}}}, } time.Sleep(10 * time.Millisecond) diff --git a/serf/config.go b/serf/config.go index 2a4fb7649..af0553c9f 100644 --- a/serf/config.go +++ b/serf/config.go @@ -5,7 +5,7 @@ package serf import ( "io" - "log" + "log/slog" "os" "time" @@ -209,7 +209,14 @@ 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 *slog.Logger + + // LogLevel is a custom log level which you provide. If LogLevel is set, it will + // use this for the internal logger. If LogLevel is not set, it will default to Info. + // You cannot specify both LogLevel and Logger at the same time, as Logger will + // override LogLevel. + // Setting LogLevel to Debug will enable debug logging for Serf. + LogLevel slog.Leveler // 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..cff6da910 100644 --- a/serf/delegate.go +++ b/serf/delegate.go @@ -5,7 +5,9 @@ package serf import ( "bytes" + "context" "fmt" + "log/slog" "github.com/armon/go-metrics" "github.com/hashicorp/go-msgpack/codec" @@ -22,7 +24,7 @@ var _ memberlist.Delegate = &delegate{} func (d *delegate) NodeMeta(limit int) []byte { roleBytes := d.serf.encodeTags(d.serf.config.Tags) if len(roleBytes) > limit { - panic(fmt.Errorf("Node tags '%v' exceeds length limit of %d bytes", d.serf.config.Tags, limit)) + panic(fmt.Errorf("node tags '%v' exceeds length limit of %d bytes", d.serf.config.Tags, limit)) } return roleBytes @@ -47,53 +49,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.LogAttrs(context.TODO(), slog.LevelError, "Error decoding leave message", slog.String("error", err.Error())) break } - d.serf.logger.Printf("[DEBUG] serf: messageLeaveType: %s", leave.Node) + d.serf.logger.LogAttrs(context.TODO(), slog.LevelDebug, "messageLeaveType", slog.String("node", 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.LogAttrs(context.TODO(), slog.LevelError, "Error decoding join message", slog.String("error", err.Error())) break } - d.serf.logger.Printf("[DEBUG] serf: messageJoinType: %s", join.Node) + d.serf.logger.LogAttrs(context.TODO(), slog.LevelDebug, "messageJoinType", slog.String("node", 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.LogAttrs(context.TODO(), slog.LevelError, "Error decoding user event message", slog.String("error", err.Error())) break } - d.serf.logger.Printf("[DEBUG] serf: messageUserEventType: %s", event.Name) + d.serf.logger.LogAttrs(context.TODO(), slog.LevelDebug, "messageUserEventType", slog.String("name", 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.LogAttrs(context.TODO(), slog.LevelError, "Error decoding query message", slog.String("error", err.Error())) break } - d.serf.logger.Printf("[DEBUG] serf: messageQueryType: %s", query.Name) + d.serf.logger.LogAttrs(context.TODO(), slog.LevelDebug, "messageQueryType", slog.String("name", 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.LogAttrs(context.TODO(), slog.LevelError, "Error decoding query response message", slog.String("error", err.Error())) break } - d.serf.logger.Printf("[DEBUG] serf: messageQueryResponseType: %v", resp.From) + d.serf.logger.LogAttrs(context.TODO(), slog.LevelDebug, "messageQueryResponseType", slog.String("from", resp.From)) d.serf.handleQueryResponse(&resp) case messageRelayType: @@ -102,7 +104,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.LogAttrs(context.TODO(), slog.LevelError, "Error decoding relay header", slog.String("error", err.Error())) break } @@ -115,14 +117,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.LogAttrs(context.TODO(), slog.LevelDebug, "Relaying response to addr", slog.String("destination", 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.LogAttrs(context.TODO(), slog.LevelError, "Error forwarding message", slog.String("destination", header.DestAddr.String()), slog.String("error", err.Error())) break } default: - d.serf.logger.Printf("[WARN] serf: Received message of unknown type: %d", t) + d.serf.logger.LogAttrs(context.TODO(), slog.LevelWarn, "Received message of unknown type", slog.Uint64("type", uint64(t))) } if rebroadcast { @@ -202,7 +204,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to encode local state", slog.String("error", err.Error())) return nil } return buf @@ -211,13 +213,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.LogAttrs(context.TODO(), slog.LevelError, "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.LogAttrs(context.TODO(), slog.LevelError, "Remote state has bad type prefix", slog.Uint64("prefix", uint64(buf[0]))) return } @@ -228,7 +230,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to decode remote state", slog.String("error", err.Error())) return } diff --git a/serf/delegate_test.go b/serf/delegate_test.go index 3174e3dde..9e7309371 100644 --- a/serf/delegate_test.go +++ b/serf/delegate_test.go @@ -174,10 +174,10 @@ func TestDelegate_MergeRemoteState(t *testing.T) { LeftMembers: []string{"foo"}, EventLTime: 50, Events: []*userEvents{ - &userEvents{ + { LTime: 45, Events: []userEvent{ - userEvent{ + { Name: "test", Payload: nil, }, diff --git a/serf/internal_query.go b/serf/internal_query.go index cac3340cf..22eb77086 100644 --- a/serf/internal_query.go +++ b/serf/internal_query.go @@ -4,9 +4,10 @@ package serf import ( + "context" "encoding/base64" "fmt" - "log" + "log/slog" "strings" ) @@ -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 *slog.Logger outCh chan<- Event serf *Serf shutdownCh <-chan struct{} @@ -75,11 +76,11 @@ 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 *slog.Logger, outCh chan<- Event, shutdownCh <-chan struct{}) (chan<- Event, error) { inCh := make(chan Event, 1024) q := &serfQueries{ inCh: inCh, - logger: logger, + logger: logger.WithGroup("serf"), outCh: outCh, serf: serf, shutdownCh: shutdownCh, @@ -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.LogAttrs(context.TODO(), slog.LevelWarn, "Unhandled internal query", slog.String("name", 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.LogAttrs(context.TODO(), slog.LevelDebug, "Got conflict resolution query", slog.String("node", 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to encode conflict query response", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to respond to conflict query: %v", slog.String("error", err.Error())) } } @@ -196,11 +197,11 @@ func (s *serfQueries) keyListResponseWithCorrectSize(q *Query, resp *nodeKeyResp } if actual > i { - s.logger.Printf("[WARN] serf: %s", resp.Message) + s.logger.LogAttrs(context.TODO(), slog.LevelWarn, resp.Message) } return raw, qresp, nil } - return nil, messageQueryResponse{}, fmt.Errorf("Failed to truncate response so that it fits into message") + return nil, messageQueryResponse{}, fmt.Errorf("failed to truncate response so that it fits into message") } // sendKeyResponse handles responding to key-related queries. @@ -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.LogAttrs(context.TODO(), slog.LevelError, err.Error()) return } if err := q.respondWithMessageAndResponse(raw, qresp); err != nil { - s.logger.Printf("[ERR] serf: Failed to respond to key query: %v", err) + s.logger.LogAttrs(context.TODO(), slog.LevelError, "Failed to respond to key query", slog.String("error", err.Error())) return } default: buf, err := encodeMessage(messageKeyResponseType, resp) if err != nil { - s.logger.Printf("[ERR] serf: Failed to encode key response: %v", err) + s.logger.LogAttrs(context.TODO(), slog.LevelError, "Failed to encode key response", slog.String("error", err.Error())) return } if err := q.Respond(buf); err != nil { - s.logger.Printf("[ERR] serf: Failed to respond to key query: %v", err) + s.logger.LogAttrs(context.TODO(), slog.LevelError, "Failed to respond to key query", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to decode key request", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "No keyring to modify (encryption not enabled)") goto SEND } - s.logger.Printf("[INFO] serf: Received install-key query") + s.logger.LogAttrs(context.TODO(), slog.LevelInfo, "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.LogAttrs(context.TODO(), slog.LevelError, "Failed to install key", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to write keyring file", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to decode key request", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "No keyring to modify (encryption not enabled)") goto SEND } - s.logger.Printf("[INFO] serf: Received use-key query") + s.logger.LogAttrs(context.TODO(), slog.LevelInfo, "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.LogAttrs(context.TODO(), slog.LevelError, "Failed to change primary key", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to write keyring file", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to decode key request", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "No keyring to modify (encryption not enabled)") goto SEND } - s.logger.Printf("[INFO] serf: Received remove-key query") + s.logger.LogAttrs(context.TODO(), slog.LevelInfo, "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.LogAttrs(context.TODO(), slog.LevelError, "Failed to remove key", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "Failed to write keyring file", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelError, "Keyring is empty (encryption not enabled)") goto SEND } - s.logger.Printf("[INFO] serf: Received list-keys query") + s.logger.LogAttrs(context.TODO(), slog.LevelInfo, "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..41a30f1ac 100644 --- a/serf/internal_query_test.go +++ b/serf/internal_query_test.go @@ -4,7 +4,7 @@ package serf import ( - "log" + "log/slog" "os" "strings" "testing" @@ -20,11 +20,15 @@ func TestInternalQueryName(t *testing.T) { func TestSerfQueries_Passthrough(t *testing.T) { serf := &Serf{} - logger := log.New(os.Stderr, "", log.LstdFlags) outCh := make(chan Event, 4) shutdown := make(chan struct{}) defer close(shutdown) - eventCh, err := newSerfQueries(serf, logger, outCh, shutdown) + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stdout, handlerOpts) + eventCh, err := newSerfQueries(serf, slog.New(handler), outCh, shutdown) if err != nil { t.Fatalf("err: %v", err) } @@ -50,11 +54,15 @@ func TestSerfQueries_Passthrough(t *testing.T) { func TestSerfQueries_Ping(t *testing.T) { serf := &Serf{} - logger := log.New(os.Stderr, "", log.LstdFlags) outCh := make(chan Event, 4) shutdown := make(chan struct{}) defer close(shutdown) - eventCh, err := newSerfQueries(serf, logger, outCh, shutdown) + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stdout, handlerOpts) + eventCh, err := newSerfQueries(serf, slog.New(handler), outCh, shutdown) if err != nil { t.Fatalf("err: %v", err) } @@ -72,11 +80,15 @@ 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) outCh := make(chan Event, 4) shutdown := make(chan struct{}) defer close(shutdown) - eventCh, err := newSerfQueries(serf, logger, outCh, shutdown) + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stdout, handlerOpts) + eventCh, err := newSerfQueries(serf, slog.New(handler), outCh, shutdown) if err != nil { t.Fatalf("err: %v", err) } @@ -124,7 +136,12 @@ func TestSerfQueries_estimateMaxKeysInListKeyResponseFactor(t *testing.T) { } func TestSerfQueries_keyListResponseWithCorrectSize(t *testing.T) { - s := serfQueries{logger: log.New(os.Stderr, "", log.LstdFlags)} + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stdout, handlerOpts) + s := serfQueries{logger: slog.New(handler)} 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..51d0294d3 100644 --- a/serf/keymanager.go +++ b/serf/keymanager.go @@ -4,8 +4,10 @@ package serf import ( + "context" "encoding/base64" "fmt" + "log/slog" "sync" ) @@ -77,7 +79,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.LogAttrs(context.TODO(), slog.LevelWarn, nodeResponse.Message) } // Currently only used for key list queries, this adds keys to a counter diff --git a/serf/merge_delegate.go b/serf/merge_delegate.go index 2c8e07c0b..f3cdbca81 100644 --- a/serf/merge_delegate.go +++ b/serf/merge_delegate.go @@ -68,11 +68,11 @@ func (m *mergeDelegate) validateMemberInfo(n *memberlist.Node) error { } if len(n.Addr) != 4 && len(n.Addr) != 16 { - return fmt.Errorf("IP byte length is invalid: %d bytes is not either 4 or 16", len(n.Addr)) + return fmt.Errorf("iP byte length is invalid: %d bytes is not either 4 or 16", len(n.Addr)) } if len(n.Meta) > memberlist.MetaMaxSize { - return fmt.Errorf("Encoded length of tags exceeds limit of %d bytes", + return fmt.Errorf("encoded length of tags exceeds limit of %d bytes", memberlist.MetaMaxSize) } return nil diff --git a/serf/ping_delegate.go b/serf/ping_delegate.go index 78a34f742..b8fb12e75 100644 --- a/serf/ping_delegate.go +++ b/serf/ping_delegate.go @@ -5,6 +5,8 @@ package serf import ( "bytes" + "context" + "log/slog" "time" "github.com/armon/go-metrics" @@ -39,7 +41,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to encode coordinate", slog.String("error", err.Error())) } return buf.Bytes() } @@ -47,14 +49,14 @@ func (p *pingDelegate) AckPayload() []byte { // NotifyPingComplete is called when this node successfully completes a direct ping // of a peer node. func (p *pingDelegate) NotifyPingComplete(other *memberlist.Node, rtt time.Duration, payload []byte) { - if payload == nil || len(payload) == 0 { + if len(payload) == 0 { return } // 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.LogAttrs(context.TODO(), slog.LevelError, "Unsupported ping version", slog.Int("version", int(version))) return } @@ -63,7 +65,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to decode coordinate from ping", slog.String("error", err.Error())) return } @@ -72,8 +74,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.LogAttrs(context.TODO(), slog.LevelDebug, "Rejected coordinate", + slog.String("from", other.Name), slog.String("error", err.Error())) return } diff --git a/serf/query.go b/serf/query.go index 2c04a9078..4ad657a5a 100644 --- a/serf/query.go +++ b/serf/query.go @@ -4,8 +4,10 @@ package serf import ( + "context" "errors" "fmt" + "log/slog" "math" "math/rand" "net" @@ -219,7 +221,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.LogAttrs(context.TODO(), slog.LevelWarn, "failed to decode filterNodeType", slog.String("error", err.Error())) return false } @@ -239,7 +241,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.LogAttrs(context.TODO(), slog.LevelWarn, "failed to decode filterTagType", slog.String("error", err.Error())) return false } @@ -247,7 +249,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.LogAttrs(context.TODO(), slog.LevelWarn, "failed to compile filter regex", slog.String("expression", filt.Expr), slog.String("error", err.Error())) return false } if !matched { @@ -255,7 +257,7 @@ func (s *Serf) shouldProcessQuery(filters [][]byte) bool { } default: - s.logger.Printf("[WARN] serf: query has unrecognized filter type: %d", filter[0]) + s.logger.LogAttrs(context.TODO(), slog.LevelWarn, "query has unrecognized filter type", slog.Uint64("type", uint64(filter[0]))) return false } } diff --git a/serf/serf.go b/serf/serf.go index 96f1d0a09..4646def63 100644 --- a/serf/serf.go +++ b/serf/serf.go @@ -5,11 +5,11 @@ package serf import ( "bytes" + "context" "encoding/base64" "encoding/json" "fmt" - "io/ioutil" - "log" + "log/slog" "math/rand" "net" "os" @@ -30,7 +30,7 @@ import ( // version to memberlist below. const ( ProtocolVersionMin uint8 = 2 - ProtocolVersionMax = 5 + ProtocolVersionMax uint8 = 5 ) const ( @@ -44,7 +44,7 @@ const MaxNodeNameLength int = 128 var ( // FeatureNotSupported is returned if a feature cannot be used // due to an older protocol version being used. - FeatureNotSupported = fmt.Errorf("Feature not supported") + FeatureNotSupported = fmt.Errorf("feature not supported") ) func init() { @@ -96,7 +96,7 @@ type Serf struct { queryResponse map[LamportTime]*QueryResponse queryLock sync.RWMutex - logger *log.Logger + logger *slog.Logger joinLock sync.Mutex stateLock sync.Mutex state SerfState @@ -249,10 +249,10 @@ const ( func Create(conf *Config) (*Serf, error) { conf.Init() if conf.ProtocolVersion < ProtocolVersionMin { - return nil, fmt.Errorf("Protocol version '%d' too low. Must be in range: [%d, %d]", + return nil, fmt.Errorf("protocol version '%d' too low. Must be in range: [%d, %d]", conf.ProtocolVersion, ProtocolVersionMin, ProtocolVersionMax) } else if conf.ProtocolVersion > ProtocolVersionMax { - return nil, fmt.Errorf("Protocol version '%d' too high. Must be in range: [%d, %d]", + return nil, fmt.Errorf("protocol version '%d' too high. Must be in range: [%d, %d]", conf.ProtocolVersion, ProtocolVersionMin, ProtocolVersionMax) } @@ -266,7 +266,13 @@ func Create(conf *Config) (*Serf, error) { if logOutput == nil { logOutput = os.Stderr } - logger = log.New(logOutput, "", log.LstdFlags) + handlerOpts := slog.HandlerOptions{} + if conf.LogLevel != nil { + handlerOpts.Level = conf.LogLevel + handlerOpts.AddSource = handlerOpts.Level == slog.LevelDebug + } + handler := slog.NewTextHandler(logOutput, &handlerOpts) + logger = slog.New(handler) } serf := &Serf{ @@ -282,7 +288,7 @@ func Create(conf *Config) (*Serf, error) { // Check that the meta data length is okay if len(serf.encodeTags(conf.Tags)) > memberlist.MetaMaxSize { - return nil, fmt.Errorf("Encoded length of tags exceeds limit of %d bytes", memberlist.MetaMaxSize) + return nil, fmt.Errorf("encoded length of tags exceeds limit of %d bytes", memberlist.MetaMaxSize) } if err := serf.ValidateNodeNames(); err != nil { return nil, err @@ -314,7 +320,7 @@ func Create(conf *Config) (*Serf, error) { // the queries outCh, err := newSerfQueries(serf, serf.logger, conf.EventCh, serf.shutdownCh) if err != nil { - return nil, fmt.Errorf("Failed to setup serf query handler: %v", err) + return nil, fmt.Errorf("failed to setup serf query handler: %v", err) } conf.EventCh = outCh @@ -324,7 +330,7 @@ func Create(conf *Config) (*Serf, error) { coordinateConfig.MetricLabels = serf.metricLabels serf.coordClient, err = coordinate.NewClient(coordinateConfig) if err != nil { - return nil, fmt.Errorf("Failed to create coordinate client: %v", err) + return nil, fmt.Errorf("failed to create coordinate client: %v", err) } } @@ -341,7 +347,7 @@ func Create(conf *Config) (*Serf, error) { conf.EventCh, serf.shutdownCh) if err != nil { - return nil, fmt.Errorf("Failed to setup snapshot: %v", err) + return nil, fmt.Errorf("failed to setup snapshot: %v", err) } snap.metricLabels = serf.metricLabels serf.snapshotter = snap @@ -420,7 +426,7 @@ func Create(conf *Config) (*Serf, error) { // and failure detection for the Serf instance. memberlist, err := memberlist.Create(conf.MemberlistConfig) if err != nil { - return nil, fmt.Errorf("Failed to create memberlist: %v", err) + return nil, fmt.Errorf("failed to create memberlist: %v", err) } serf.memberlist = memberlist @@ -549,7 +555,7 @@ func (s *Serf) Query(name string, payload []byte, params *QueryParam) (*QueryRes // Encode the filters filters, err := params.encodeFilters() if err != nil { - return nil, fmt.Errorf("Failed to format filters: %v", err) + return nil, fmt.Errorf("failed to format filters: %v", err) } // Setup the flags @@ -623,7 +629,7 @@ func (s *Serf) registerQueryResponse(timeout time.Duration, resp *QueryResponse) func (s *Serf) SetTags(tags map[string]string) error { // Check that the meta data length is okay if len(s.encodeTags(tags)) > memberlist.MetaMaxSize { - return fmt.Errorf("Encoded length of tags exceeds limit of %d bytes", + return fmt.Errorf("encoded length of tags exceeds limit of %d bytes", memberlist.MetaMaxSize) } @@ -641,7 +647,7 @@ func (s *Serf) SetTags(tags map[string]string) error { func (s *Serf) Join(existing []string, ignoreOld bool) (int, error) { // Do a quick state check if s.State() != SerfAlive { - return 0, fmt.Errorf("Serf can't Join after Leave or Shutdown") + return 0, fmt.Errorf("serf can't Join after Leave or Shutdown") } // Hold the joinLock, this is to make eventJoinIgnore safe @@ -688,7 +694,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.LogAttrs(context.TODO(), slog.LevelWarn, "Failed to broadcast join intent", slog.String("error", err.Error())) return err } return nil @@ -705,10 +711,10 @@ func (s *Serf) Leave() error { return nil } else if s.state == SerfLeaving { s.stateLock.Unlock() - return fmt.Errorf("Leave already in progress") + return fmt.Errorf("leave already in progress") } else if s.state == SerfShutdown { s.stateLock.Unlock() - return fmt.Errorf("Leave called after Shutdown") + return fmt.Errorf("leave called after Shutdown") } s.state = SerfLeaving s.stateLock.Unlock() @@ -739,14 +745,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.LogAttrs(context.TODO(), slog.LevelWarn, "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.LogAttrs(context.TODO(), slog.LevelWarn, "timeout waiting for leave broadcast", slog.String("error", err.Error())) } // Wait for the leave to propagate through the cluster. The broadcast @@ -870,7 +876,7 @@ func (s *Serf) Shutdown() error { } if s.state != SerfLeft { - s.logger.Printf("[WARN] serf: Shutdown without a Leave") + s.logger.LogAttrs(context.TODO(), slog.LevelWarn, "Shutdown without a Leave") } // Wait to close the shutdown channel until after we've shut down the @@ -995,8 +1001,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.LogAttrs(context.TODO(), slog.LevelInfo, "EventMemberJoin", + slog.String("name", member.Member.Name), slog.String("address", member.Member.Addr.String())) if s.config.EventCh != nil { s.config.EventCh <- MemberEvent{ Type: EventMemberJoin, @@ -1029,7 +1035,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.LogAttrs(context.TODO(), slog.LevelWarn, "Bad state when leave", slog.String("status", member.Status.String())) return } @@ -1044,8 +1050,7 @@ 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.LogAttrs(context.TODO(), slog.LevelInfo, eventStr, slog.String("name", member.Member.Name), slog.String("address", member.Member.Addr.String())) if s.config.EventCh != nil { s.config.EventCh <- MemberEvent{ Type: event, @@ -1089,7 +1094,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.LogAttrs(context.TODO(), slog.LevelInfo, "EventMemberUpdate", slog.String("name", member.Member.Name)) if s.config.EventCh != nil { s.config.EventCh <- MemberEvent{ Type: EventMemberUpdate, @@ -1122,7 +1127,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.LogAttrs(context.TODO(), slog.LevelDebug, "Refuting an older leave intent") go s.broadcastJoin(s.clock.Time()) return false } @@ -1166,8 +1171,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.LogAttrs(context.TODO(), slog.LevelInfo, "[EventMemberLeave (forced)", + slog.String("name", member.Member.Name), slog.String("address", member.Member.Addr.String())) if s.config.EventCh != nil { s.config.EventCh <- MemberEvent{ Type: EventMemberLeave, @@ -1198,7 +1203,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.LogAttrs(context.TODO(), slog.LevelInfo, "EventMemberReap (forced)", slog.String("name", member.Member.Name), slog.String("address", member.Member.Addr.String())) //If we are leaving or left we may be in that list of members if member.Status == StatusLeaving || member.Status == StatusLeft { @@ -1257,11 +1262,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.LogAttrs(context.TODO(), slog.LevelWarn, + "received old event", + slog.String("event", eventMsg.Name), + slog.Uint64("time", uint64(eventMsg.LTime)), + slog.Uint64("current time", uint64(s.eventClock.Time()))) return false } @@ -1316,11 +1321,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.LogAttrs(context.TODO(), slog.LevelWarn, + "received old query", + slog.String("event", query.Name), + slog.Uint64("time", uint64(query.LTime)), + slog.Uint64("current time", uint64(s.queryClock.Time()))) return false } @@ -1369,7 +1374,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.LogAttrs(context.TODO(), slog.LevelError, "failed to format ack", slog.String("error", err.Error())) } else { udpAddr := net.UDPAddr{IP: query.Addr, Port: int(query.Port)} addr := memberlist.Address{ @@ -1377,10 +1382,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.LogAttrs(context.TODO(), slog.LevelError, "failed to send ack", slog.String("error", err.Error())) } 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.LogAttrs(context.TODO(), slog.LevelError, "failed to relay ack", slog.String("error", err.Error())) } } } @@ -1410,15 +1415,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.LogAttrs(context.TODO(), slog.LevelWarn, "reply for non-running query (LTime: %d, ID: %d) From: %s", + slog.Uint64("lamport time", uint64(resp.LTime)), slog.Uint64("ID", uint64(resp.ID)), slog.String("from", 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.LogAttrs(context.TODO(), slog.LevelWarn, "query reply ID mismatch (Local: %d, Response: %d)", + slog.Uint64("local", uint64(query.id)), slog.Uint64("response", uint64(resp.ID))) return } @@ -1438,7 +1443,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.LogAttrs(context.TODO(), slog.LevelWarn, err.Error()) } } else { // Exit early if this is a duplicate response @@ -1450,7 +1455,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.LogAttrs(context.TODO(), slog.LevelWarn, err.Error()) } } } @@ -1461,14 +1466,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.LogAttrs(context.TODO(), slog.LevelWarn, "Name conflict", + slog.String("name", existing.Name), slog.Group("existing", slog.String("addr", existing.Addr.String()), slog.Uint64("port", uint64(existing.Port))), slog.Group("other", slog.String("addr", other.Addr.String()), slog.Uint64("port", uint64(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.LogAttrs(context.TODO(), slog.LevelError, "Node name conflicts with another node. Names must be unique! (Resolution enabled: %v)", + slog.Group("conflicting node", slog.String("address", other.Addr.String()), slog.Uint64("port", uint64(other.Port))), slog.Bool("Resolution enabled", s.config.EnableNameConflictResolution)) // If automatic resolution is enabled, kick off the resolution if s.config.EnableNameConflictResolution { @@ -1487,7 +1492,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to start name resolution query", slog.String("error", err.Error())) return } @@ -1499,12 +1504,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.LogAttrs(context.TODO(), slog.LevelError, "Invalid conflict query response type", slog.String("payload", string(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.LogAttrs(context.TODO(), slog.LevelError, "Failed to decode conflict query response", slog.String("error", err.Error())) continue } @@ -1518,20 +1523,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.LogAttrs(context.TODO(), slog.LevelInfo, "majority in name conflict resolution", + slog.Int("matching", matching), slog.Int("responses", responses)) return } // Since we lost the vote, we need to exit - s.logger.Printf("[WARN] serf: minority in name conflict resolution, quiting [%d / %d]", - matching, responses) + s.logger.LogAttrs(context.TODO(), slog.LevelWarn, "minority in name conflict resolution, quiting", + slog.Int("matching", matching), slog.Int("responses", responses)) if err := s.Shutdown(); err != nil { - s.logger.Printf("[ERR] serf: Failed to shutdown: %v", err) + s.logger.LogAttrs(context.TODO(), slog.LevelError, "Failed to shutdown", slog.String("error", err.Error())) } } -//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 +1616,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.LogAttrs(context.TODO(), slog.LevelInfo, "EventMemberReap", slog.String("name", m.Name)) s.eraseNode(m) } @@ -1643,7 +1648,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.LogAttrs(context.TODO(), slog.LevelDebug, "forgoing reconnect for random throttling") return } @@ -1653,7 +1658,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.LogAttrs(context.TODO(), slog.LevelInfo, "attempting reconnect", slog.String("name", mem.Name), slog.String("address", addr.String())) joinAddr := addr.String() if mem.Name != "" { @@ -1690,11 +1695,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.LogAttrs(context.TODO(), slog.LevelWarn, "%s queue depth: %d", slog.String("name", name), slog.Int("amount queued", 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.LogAttrs(context.TODO(), slog.LevelWarn, "%s queue depth (%d) exceeds limit (%d), dropping messages!", + slog.String("name", name), slog.Int("queue depth", numq), slog.Int("limit", max)) queue.Prune(max) } case <-s.shutdownCh: @@ -1772,14 +1777,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.LogAttrs(context.TODO(), slog.LevelInfo, "Attempting re-join to previously known node", slog.String("previous", prev.String())) _, err := s.memberlist.Join([]string{joinAddr}) if err == nil { - s.logger.Printf("[INFO] serf: Re-joined to previously known node: %s", prev) + s.logger.LogAttrs(context.TODO(), slog.LevelInfo, "Re-joined to previously known node", slog.String("previous", prev.String())) return } } - s.logger.Printf("[WARN] serf: Failed to re-join any previously known node") + s.logger.LogAttrs(context.TODO(), slog.LevelWarn, "Failed to re-join any previously known node") } // encodeTags is used to encode a tag map @@ -1814,7 +1819,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to decode tags", slog.String("error", err.Error())) } return tags } @@ -1866,12 +1871,12 @@ func (s *Serf) writeKeyringFile() error { encodedKeys, err := json.MarshalIndent(keysEncoded, "", " ") if err != nil { - return fmt.Errorf("Failed to encode keys: %s", err) + return fmt.Errorf("failed to encode keys: %s", err) } // Use 0600 for permissions because key data is sensitive - if err = ioutil.WriteFile(s.config.KeyringFile, encodedKeys, 0600); err != nil { - return fmt.Errorf("Failed to write keyring file: %s", err) + if err = os.WriteFile(s.config.KeyringFile, encodedKeys, 0600); err != nil { + return fmt.Errorf("failed to write keyring file: %s", err) } // Success! @@ -1884,7 +1889,7 @@ func (s *Serf) GetCoordinate() (*coordinate.Coordinate, error) { return s.coordClient.GetCoordinate(), nil } - return nil, fmt.Errorf("Coordinates are disabled") + return nil, fmt.Errorf("coordinates are disabled") } // GetCachedCoordinate returns the network coordinate for the node with the given @@ -1923,11 +1928,11 @@ func (s *Serf) validateNodeName(name string) error { if s.config.ValidateNodeNames { var InvalidNameRe = regexp.MustCompile(`[^A-Za-z0-9\-\.]+`) if InvalidNameRe.MatchString(name) { - return fmt.Errorf("Node name contains invalid characters %v , Valid characters include "+ + return fmt.Errorf("node name contains invalid characters %v , Valid characters include "+ "all alpha-numerics and dashes and '.' ", name) } if len(name) > MaxNodeNameLength { - return fmt.Errorf("Node name is %v characters. "+ + return fmt.Errorf("node name is %v characters. "+ "Valid length is between 1 and 128 characters", len(name)) } } diff --git a/serf/serf_test.go b/serf/serf_test.go index 0dada60fb..4e8f55abd 100644 --- a/serf/serf_test.go +++ b/serf/serf_test.go @@ -9,7 +9,7 @@ import ( "encoding/base64" "fmt" "io/ioutil" - "log" + "log/slog" "net" "os" "path/filepath" @@ -59,8 +59,13 @@ 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 + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stdout, handlerOpts) + config.Logger = slog.New(handler) + config.MemberlistConfig.Logger = slog.NewLogLogger(handler, slog.LevelDebug) } return config @@ -1116,7 +1121,7 @@ func TestSerf_update(t *testing.T) { t.Fatalf("err: %v", err) } - if time.Now().Sub(start) > 2*time.Second { + if time.Since(start) > 2*time.Second { t.Fatalf("timed out trying to restart") } } @@ -1497,9 +1502,9 @@ func TestSerf_Reap(t *testing.T) { m := Member{} old := []*memberState{ - &memberState{m, 0, time.Now()}, - &memberState{m, 0, time.Now().Add(-5 * time.Second)}, - &memberState{m, 0, time.Now().Add(-10 * time.Second)}, + {m, 0, time.Now()}, + {m, 0, time.Now().Add(-5 * time.Second)}, + {m, 0, time.Now().Add(-10 * time.Second)}, } old = s.reap(old, time.Now(), time.Second*6) @@ -1510,9 +1515,9 @@ func TestSerf_Reap(t *testing.T) { func TestRemoveOldMember(t *testing.T) { old := []*memberState{ - &memberState{Member: Member{Name: "foo"}}, - &memberState{Member: Member{Name: "bar"}}, - &memberState{Member: Member{Name: "baz"}}, + {Member: Member{Name: "foo"}}, + {Member: Member{Name: "bar"}}, + {Member: Member{Name: "baz"}}, } old = removeOldMember(old, "bar") @@ -1843,7 +1848,7 @@ func TestSerf_SnapshotRecovery(t *testing.T) { // Wait for the node to auto rejoin start := time.Now() - for time.Now().Sub(start) < time.Second { + for time.Since(start) < time.Second { members := s1.Members() if len(members) == 2 && members[0].Status == StatusAlive && members[1].Status == StatusAlive { break @@ -2595,7 +2600,7 @@ type CancelMergeDelegate struct { func (c *CancelMergeDelegate) NotifyMerge(members []*Member) error { c.invoked = true - return fmt.Errorf("Merge canceled") + return fmt.Errorf("merge canceled") } func TestSerf_Join_Cancel(t *testing.T) { diff --git a/serf/snapshot.go b/serf/snapshot.go index 28bb668e6..36c0dad99 100644 --- a/serf/snapshot.go +++ b/serf/snapshot.go @@ -5,8 +5,9 @@ package serf import ( "bufio" + "context" "fmt" - "log" + "log/slog" "math/rand" "net" "os" @@ -72,7 +73,7 @@ type Snapshotter struct { lastQueryClock LamportTime leaveCh chan struct{} leaving bool - logger *log.Logger + logger *slog.Logger minCompactSize int64 path string offset int64 @@ -103,7 +104,7 @@ func (p PreviousNode) String() string { func NewSnapshotter(path string, minCompactSize int, rejoinAfterLeave bool, - logger *log.Logger, + logger *slog.Logger, clock *LamportClock, outCh chan<- Event, shutdownCh <-chan struct{}) (chan<- Event, *Snapshotter, error) { @@ -262,7 +263,7 @@ func (s *Snapshotter) stream() { case *Query: s.processQuery(typed) default: - s.logger.Printf("[ERR] serf: Unknown event to snapshot: %#v", e) + s.logger.LogAttrs(context.TODO(), slog.LevelError, "Unknown event to snapshot", slog.String("error", e.String())) } } @@ -277,10 +278,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.LogAttrs(context.TODO(), slog.LevelError, "failed to flush leave to snapshot", slog.String("error", err.Error())) } if err := s.fh.Sync(); err != nil { - s.logger.Printf("[ERR] serf: failed to sync leave to snapshot: %v", err) + s.logger.LogAttrs(context.TODO(), slog.LevelError, "failed to sync leave to snapshot", slog.String("error", err.Error())) } case e := <-s.streamCh: @@ -310,10 +311,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.LogAttrs(context.TODO(), slog.LevelError, "failed to flush snapshot", slog.String("error", err.Error())) } if err := s.fh.Sync(); err != nil { - s.logger.Printf("[ERR] serf: failed to sync snapshot: %v", err) + s.logger.LogAttrs(context.TODO(), slog.LevelError, "failed to sync snapshot", slog.String("error", err.Error())) } s.fh.Close() close(s.waitCh) @@ -377,16 +378,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.LogAttrs(context.TODO(), slog.LevelError, "Failed to update snapshot", slog.String("error", err.Error())) 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.LogAttrs(context.TODO(), slog.LevelInfo, "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.LogAttrs(context.TODO(), slog.LevelError, "Compaction failed, will reattempt after snapshot recovery interval", slog.Duration("interval", snapshotErrorRecoveryInterval), slog.String("error", err.Error())) } else { - s.logger.Printf("[INFO] serf: Finished compaction, successfully recovered from error state") + s.logger.LogAttrs(context.TODO(), slog.LevelInfo, "Finished compaction, successfully recovered from error state") } } } @@ -564,7 +565,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.LogAttrs(context.TODO(), slog.LevelWarn, "Failed to parse address", slog.String("address", line)) continue } addr := info[addrIdx+1:] @@ -579,7 +580,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.LogAttrs(context.TODO(), slog.LevelWarn, "Failed to convert clock time", slog.String("error", err.Error())) continue } s.lastClock = LamportTime(timeInt) @@ -588,7 +589,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.LogAttrs(context.TODO(), slog.LevelWarn, "Failed to convert event clock time", slog.String("error", err.Error())) continue } s.lastEventClock = LamportTime(timeInt) @@ -597,7 +598,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.LogAttrs(context.TODO(), slog.LevelWarn, "Failed to convert query clock time", slog.String("error", err.Error())) continue } s.lastQueryClock = LamportTime(timeInt) @@ -607,7 +608,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.LogAttrs(context.TODO(), slog.LevelInfo, "Ignoring previous leave in snapshot") continue } s.aliveNodes = make(map[string]string) @@ -619,7 +620,7 @@ func (s *Snapshotter) replay() error { // Skip comment lines } else { - s.logger.Printf("[WARN] serf: Unrecognized snapshot line: %v", line) + s.logger.LogAttrs(context.TODO(), slog.LevelWarn, "Unrecognized snapshot line", slog.String("line", line)) } } diff --git a/serf/snapshot_test.go b/serf/snapshot_test.go index bab5bc76e..5e11bf939 100644 --- a/serf/snapshot_test.go +++ b/serf/snapshot_test.go @@ -6,8 +6,10 @@ package serf import ( "fmt" "io/ioutil" - "log" + "log/slog" + "math/rand" "os" + "path" "reflect" "testing" "time" @@ -23,7 +25,8 @@ func TestSnapshotter(t *testing.T) { clock := new(LamportClock) outCh := make(chan Event, 64) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + handler := slog.NewTextHandler(os.Stdout, &slog.HandlerOptions{}) + logger := slog.New(handler) inCh, snap, err := NewSnapshotter(td+"snap", snapshotSizeLimit, false, logger, clock, outCh, stopCh) if err != nil { @@ -168,19 +171,20 @@ func TestSnapshotter(t *testing.T) { } func TestSnapshotter_forceCompact(t *testing.T) { - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) clock := new(LamportClock) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stderr, handlerOpts) + logger := slog.New(handler) // Create a very low limit - inCh, snap, err := NewSnapshotter(td+"snap", 1024, false, - logger, clock, nil, stopCh) + inCh, snap, err := NewSnapshotter(td+"snap", 1024, false, logger, clock, nil, stopCh) if err != nil { t.Fatalf("err: %v", err) } @@ -232,15 +236,18 @@ func TestSnapshotter_forceCompact(t *testing.T) { } func TestSnapshotter_leave(t *testing.T) { - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) clock := new(LamportClock) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stderr, handlerOpts) + logger := slog.New(handler) + inCh, snap, err := NewSnapshotter(td+"snap", snapshotSizeLimit, false, logger, clock, nil, stopCh) if err != nil { @@ -313,15 +320,18 @@ func TestSnapshotter_leave(t *testing.T) { } func TestSnapshotter_leave_rejoin(t *testing.T) { - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) clock := new(LamportClock) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stderr, handlerOpts) + logger := slog.New(handler) + inCh, snap, err := NewSnapshotter(td+"snap", snapshotSizeLimit, true, logger, clock, nil, stopCh) if err != nil { @@ -395,16 +405,17 @@ func TestSnapshotter_leave_rejoin(t *testing.T) { func TestSnapshotter_slowDiskNotBlockingEventCh(t *testing.T) { t.Skip("Flaky test") - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } - t.Log("Temp dir", td) + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) clock := new(LamportClock) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stderr, handlerOpts) + logger := slog.New(handler) outCh := make(chan Event, 1024) inCh, snap, err := NewSnapshotter(td+"snap", snapshotSizeLimit, true, @@ -481,16 +492,17 @@ func TestSnapshotter_slowDiskNotBlockingEventCh(t *testing.T) { } func TestSnapshotter_blockedUpstreamNotBlockingMemberlist(t *testing.T) { - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } - t.Log("Temp dir", td) + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) clock := new(LamportClock) stopCh := make(chan struct{}) - logger := log.New(os.Stderr, "", log.LstdFlags) + handlerOpts := &slog.HandlerOptions{ + AddSource: true, + Level: slog.LevelDebug, + } + handler := slog.NewTextHandler(os.Stderr, handlerOpts) + logger := slog.New(handler) // 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 ( From 726ad91571cd38ca6e970a29009f92ac2472c3b9 Mon Sep 17 00:00:00 2001 From: Luke Thorne Date: Sun, 20 Aug 2023 18:07:02 -0500 Subject: [PATCH 2/3] remove ioutil --- cmd/serf/command/agent/agent_test.go | 29 +++++++------------- cmd/serf/command/agent/config_test.go | 17 ++++++------ cmd/serf/command/agent/event_handler_test.go | 11 ++++---- cmd/serf/main.go | 4 +-- serf/serf_test.go | 24 ++++++---------- serf/snapshot_test.go | 6 +--- 6 files changed, 34 insertions(+), 57 deletions(-) diff --git a/cmd/serf/command/agent/agent_test.go b/cmd/serf/command/agent/agent_test.go index 938cebd5d..df95dcd14 100644 --- a/cmd/serf/command/agent/agent_test.go +++ b/cmd/serf/command/agent/agent_test.go @@ -5,8 +5,10 @@ package agent import ( "encoding/json" - "io/ioutil" + "fmt" + "math/rand" "os" + "path" "path/filepath" "reflect" "strings" @@ -128,10 +130,7 @@ func TestAgentTagsFile(t *testing.T) { "datacenter": "us-east", } - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) ip1, returnFn1 := testutil.TakeIP() @@ -153,9 +152,7 @@ func TestAgentTagsFile(t *testing.T) { testutil.Yield() - err = a1.SetTags(tags) - - if err != nil { + if err := a1.SetTags(tags); err != nil { t.Fatalf("err: %v", err) } @@ -248,10 +245,7 @@ func TestAgentKeyringFile(t *testing.T) { "5K9OtfP7efFrNKe5WCQvXvnaXJ5cWP0SvXiwe0kkjM4=", } - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) keyringFile := filepath.Join(td, "keyring.json") @@ -265,7 +259,7 @@ func TestAgentKeyringFile(t *testing.T) { t.Fatalf("err: %v", err) } - if err := ioutil.WriteFile(keyringFile, encodedKeys, 0600); err != nil { + if err := os.WriteFile(keyringFile, encodedKeys, 0600); err != nil { t.Fatalf("err: %v", err) } @@ -299,21 +293,18 @@ func TestAgentKeyringFile_BadOptions(t *testing.T) { } func TestAgentKeyringFile_NoKeys(t *testing.T) { - dir, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + dir := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(dir) keysFile := filepath.Join(dir, "keyring") - if err := ioutil.WriteFile(keysFile, []byte("[]"), 0600); err != nil { + if err := os.WriteFile(keysFile, []byte("[]"), 0600); err != nil { t.Fatalf("err: %v", err) } agentConfig := DefaultConfig() agentConfig.KeyringFile = keysFile - _, err = Create(agentConfig, serf.DefaultConfig(), nil) + _, err := Create(agentConfig, serf.DefaultConfig(), nil) if err == nil { t.Fatalf("should have errored") } diff --git a/cmd/serf/command/agent/config_test.go b/cmd/serf/command/agent/config_test.go index add6717ce..946f52ed4 100644 --- a/cmd/serf/command/agent/config_test.go +++ b/cmd/serf/command/agent/config_test.go @@ -6,8 +6,10 @@ package agent import ( "bytes" "encoding/base64" - "io/ioutil" + "fmt" + "math/rand" "os" + "path" "path/filepath" "reflect" "testing" @@ -535,7 +537,7 @@ func TestReadConfigPaths_badPath(t *testing.T) { } func TestReadConfigPaths_file(t *testing.T) { - tf, err := ioutil.TempFile("", "serf") + tf, err := os.CreateTemp("", "serf") if err != nil { t.Fatalf("err: %v", err) } @@ -554,26 +556,23 @@ func TestReadConfigPaths_file(t *testing.T) { } func TestReadConfigPaths_dir(t *testing.T) { - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) - err = ioutil.WriteFile(filepath.Join(td, "a.json"), + err := os.WriteFile(filepath.Join(td, "a.json"), []byte(`{"node_name": "bar"}`), 0644) if err != nil { t.Fatalf("err: %v", err) } - err = ioutil.WriteFile(filepath.Join(td, "b.json"), + err = os.WriteFile(filepath.Join(td, "b.json"), []byte(`{"node_name": "baz"}`), 0644) if err != nil { t.Fatalf("err: %v", err) } // A non-json file, shouldn't be read - err = ioutil.WriteFile(filepath.Join(td, "c"), + err = os.WriteFile(filepath.Join(td, "c"), []byte(`{"node_name": "bad"}`), 0644) if err != nil { t.Fatalf("err: %v", err) diff --git a/cmd/serf/command/agent/event_handler_test.go b/cmd/serf/command/agent/event_handler_test.go index abbe954d1..e13b53f0e 100644 --- a/cmd/serf/command/agent/event_handler_test.go +++ b/cmd/serf/command/agent/event_handler_test.go @@ -5,7 +5,6 @@ package agent import ( "fmt" - "io/ioutil" "net" "os" "testing" @@ -51,7 +50,7 @@ done // agent. It returns the path to the event script itself and a path to // the file that will contain the events that that script receives. func testEventScript(t *testing.T, script string) (string, string) { - scriptFile, err := ioutil.TempFile("", "serf") + scriptFile, err := os.CreateTemp("", "serf") if err != nil { t.Fatalf("err: %v", err) } @@ -61,7 +60,7 @@ func testEventScript(t *testing.T, script string) (string, string) { t.Fatalf("err: %v", err) } - resultFile, err := ioutil.TempFile("", "serf-result") + resultFile, err := os.CreateTemp("", "serf-result") if err != nil { t.Fatalf("err: %v", err) } @@ -111,7 +110,7 @@ func TestScriptEventHandler(t *testing.T) { h.HandleEvent(event) - result, err := ioutil.ReadFile(results) + result, err := os.ReadFile(results) if err != nil { t.Fatalf("err: %v", err) } @@ -152,7 +151,7 @@ func TestScriptUserEventHandler(t *testing.T) { h.HandleEvent(userEvent) - result, err := ioutil.ReadFile(results) + result, err := os.ReadFile(results) if err != nil { t.Fatalf("err: %v", err) } @@ -191,7 +190,7 @@ func TestScriptQueryEventHandler(t *testing.T) { h.HandleEvent(query) - result, err := ioutil.ReadFile(results) + result, err := os.ReadFile(results) if err != nil { t.Fatalf("err: %v", err) } diff --git a/cmd/serf/main.go b/cmd/serf/main.go index 69cafd5ec..7140d5351 100644 --- a/cmd/serf/main.go +++ b/cmd/serf/main.go @@ -5,7 +5,7 @@ package main import ( "fmt" - "io/ioutil" + "io" "log" "os" @@ -13,7 +13,7 @@ import ( ) func main() { - log.SetOutput(ioutil.Discard) + log.SetOutput(io.Discard) // Get the command line args. We shortcut "--version" and "-v" to // just show the version. diff --git a/serf/serf_test.go b/serf/serf_test.go index 4e8f55abd..80a085952 100644 --- a/serf/serf_test.go +++ b/serf/serf_test.go @@ -8,10 +8,11 @@ import ( "context" "encoding/base64" "fmt" - "io/ioutil" "log/slog" + "math/rand" "net" "os" + "path" "path/filepath" "reflect" "strconv" @@ -1772,10 +1773,7 @@ func TestSerf_Join_IgnoreOld(t *testing.T) { } func TestSerf_SnapshotRecovery(t *testing.T) { - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) ip1, returnFn1 := testutil.TakeIP() @@ -1869,10 +1867,7 @@ func TestSerf_Leave_SnapshotRecovery(t *testing.T) { t.Skip("test contains a data race") } - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) ip1, returnFn1 := testutil.TakeIP() @@ -2450,10 +2445,7 @@ func TestSerf_WriteKeyringFile(t *testing.T) { existing := "T9jncgl9mbLus+baTTa7q7nPSUrXwbDi2dhbtqir37s=" newKey := "HvY8ubRZMgafUOWvrOadwOckVa1wN3QWAo46FVKbVN8=" - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) keyringFile := filepath.Join(td, "tags.json") @@ -2487,7 +2479,7 @@ func TestSerf_WriteKeyringFile(t *testing.T) { t.Fatalf("err: %v", err) } - content, err := ioutil.ReadFile(keyringFile) + content, err := os.ReadFile(keyringFile) if err != nil { t.Fatalf("err: %v", err) } @@ -2517,7 +2509,7 @@ func TestSerf_WriteKeyringFile(t *testing.T) { t.Fatalf("err: %v", err) } - content, err = ioutil.ReadFile(keyringFile) + content, err = os.ReadFile(keyringFile) if err != nil { t.Fatalf("err: %v", err) } @@ -2537,7 +2529,7 @@ func TestSerf_WriteKeyringFile(t *testing.T) { t.Fatalf("err: %v", err) } - content, err = ioutil.ReadFile(keyringFile) + content, err = os.ReadFile(keyringFile) if err != nil { t.Fatalf("err: %v", err) } diff --git a/serf/snapshot_test.go b/serf/snapshot_test.go index 5e11bf939..ef0deef93 100644 --- a/serf/snapshot_test.go +++ b/serf/snapshot_test.go @@ -5,7 +5,6 @@ package serf import ( "fmt" - "io/ioutil" "log/slog" "math/rand" "os" @@ -16,10 +15,7 @@ import ( ) func TestSnapshotter(t *testing.T) { - td, err := ioutil.TempDir("", "serf") - if err != nil { - t.Fatalf("err: %v", err) - } + td := path.Join(os.TempDir(), fmt.Sprintf("serf-%d", rand.Int())) defer os.RemoveAll(td) clock := new(LamportClock) From 3fcbe0e92ec04cb346487c93403d3e466d4cc567 Mon Sep 17 00:00:00 2001 From: Luke Thorne Date: Sun, 20 Aug 2023 18:10:30 -0500 Subject: [PATCH 3/3] minor updates to match go formats --- serf/serf.go | 15 +++++---------- 1 file changed, 5 insertions(+), 10 deletions(-) diff --git a/serf/serf.go b/serf/serf.go index 4646def63..730665238 100644 --- a/serf/serf.go +++ b/serf/serf.go @@ -42,16 +42,11 @@ const ( const MaxNodeNameLength int = 128 var ( - // FeatureNotSupported is returned if a feature cannot be used + // ErrFeatureNotSupported is returned if a feature cannot be used // due to an older protocol version being used. - FeatureNotSupported = fmt.Errorf("feature not supported") + ErrFeatureNotSupported = fmt.Errorf("feature not supported") ) -func init() { - // Seed the random number generator - rand.Seed(time.Now().UnixNano()) -} - // ReconnectTimeoutOverrider is an interface that can be implemented to allow overriding // the reconnect timeout for individual members. type ReconnectTimeoutOverrider interface { @@ -218,7 +213,7 @@ func (ue *userEvent) Equals(other *userEvent) bool { if ue.Name != other.Name { return false } - if bytes.Compare(ue.Payload, other.Payload) != 0 { + if !bytes.Equal(ue.Payload, other.Payload) { return false } return true @@ -539,7 +534,7 @@ func (s *Serf) UserEvent(name string, payload []byte, coalesce bool) error { func (s *Serf) Query(name string, payload []byte, params *QueryParam) (*QueryResponse, error) { // Check that the latest protocol is in use if s.ProtocolVersion() < 4 { - return nil, FeatureNotSupported + return nil, ErrFeatureNotSupported } // Provide default parameters if none given @@ -969,7 +964,7 @@ func (s *Serf) handleNodeJoin(n *memberlist.Node) { s.members[n.Name] = member } else { oldStatus = member.Status - deadTime := time.Now().Sub(member.leaveTime) + deadTime := time.Since(member.leaveTime) if oldStatus == StatusFailed && deadTime < s.config.FlapTimeout { metrics.IncrCounterWithLabels([]string{"serf", "member", "flap"}, 1, s.metricLabels) }