diff --git a/dax/computer/service/computer.go b/dax/computer/service/computer.go index f60462414..32a2aae28 100644 --- a/dax/computer/service/computer.go +++ b/dax/computer/service/computer.go @@ -35,7 +35,7 @@ func New(addr dax.Address, cfg CommandConfig, logger logger.Logger) *computerSer return &computerService{ addr: addr, cfg: cfg, - logger: logger, + logger: logger.WithPrefix("Computer: "), } } diff --git a/dax/controller/controller.go b/dax/controller/controller.go index 0fbf0f520..5be60d9aa 100644 --- a/dax/controller/controller.go +++ b/dax/controller/controller.go @@ -108,7 +108,7 @@ func (c *Controller) Start() error { c.poller.Run() go c.nodeRegistrationRoutine(c.nodeChan, c.registrationBatchTimeout) - go c.snappingTurtleRoutine(c.snappingTurtleTimeout, c.snapControl) + go c.snappingTurtleRoutine(c.snappingTurtleTimeout, c.snapControl, c.logger.WithPrefix("Snapping Turtle: ")) return nil } diff --git a/dax/controller/node_registerer.go b/dax/controller/node_registerer.go index 7f655a643..4c4ae46ac 100644 --- a/dax/controller/node_registerer.go +++ b/dax/controller/node_registerer.go @@ -43,10 +43,9 @@ func (c *Controller) nodeRegistrationDelayed(nodes chan *dax.Node, timeout time. case <-c.stopping: return nil case node := <-nodes: - c.logger.Debugf("adding node: %+v", node) + c.logger.Printf("Adding node to registration batch: %+v", node) batch = append(batch, node) case <-time.After(timeout): - c.logger.Debugf("no new nodes in last %s, batch: %d", timeout, len(batch)) if len(batch) > 0 { err := c.RegisterNodes(context.Background(), batch...) if err != nil { diff --git a/dax/controller/service/controller.go b/dax/controller/service/controller.go index 027c6d5d9..56354888f 100644 --- a/dax/controller/service/controller.go +++ b/dax/controller/service/controller.go @@ -35,7 +35,7 @@ func New(uri *fbnet.URI, cfg controller.Config) *controllerService { // Set up logger. var logr logger.Logger = logger.StderrLogger if cfg.Logger != nil { - logr = cfg.Logger + logr = cfg.Logger.WithPrefix("Controller: ") } // Storage methods. diff --git a/dax/controller/snapping_turtle.go b/dax/controller/snapping_turtle.go index 76ca5eb84..c09a0bbda 100644 --- a/dax/controller/snapping_turtle.go +++ b/dax/controller/snapping_turtle.go @@ -5,9 +5,10 @@ import ( "time" "github.com/featurebasedb/featurebase/v3/dax" + "github.com/featurebasedb/featurebase/v3/logger" ) -func (c *Controller) snappingTurtleRoutine(period time.Duration, control chan struct{}) { +func (c *Controller) snappingTurtleRoutine(period time.Duration, control chan struct{}, log logger.Logger) { if period == 0 { return // disable automatic snapshotting } @@ -16,45 +17,45 @@ func (c *Controller) snappingTurtleRoutine(period time.Duration, control chan st select { case <-c.stopping: ticker.Stop() - c.logger.Debugf("TURTLE: Stopping Snapping Turtle") + log.Debugf("Stopping Snapping Turtle") return case <-ticker.C: - c.snapAll() + c.snapAll(log) case <-control: - c.snapAll() + c.snapAll(log) } } - } -func (c *Controller) snapAll() { - c.logger.Debugf("TURTLE: snapAll") +func (c *Controller) snapAll(log logger.Logger) { + start := time.Now() + defer func() { + log.Printf("full snapshot took: %v", time.Since(start)) + }() ctx := context.Background() tx, err := c.BoltDB.BeginTx(ctx, false) if err != nil { - c.logger.Printf("Error getting transaction for snapping turtle: %v", err) + log.Printf("Error getting transaction for snapping turtle: %v", err) return } defer tx.Rollback() qdbs, err := c.Schemar.Databases(tx, "") if err != nil { - c.logger.Printf("couldn't get databases: %v", err) + log.Printf("couldn't get databases: %v", err) } for _, qdb := range qdbs { - c.snapAllForDatabase(tx, qdb.QualifiedID()) + c.snapAllForDatabase(tx, qdb.QualifiedID(), log) } - - c.logger.Debugf("TURTLE: snapAll complete") } -func (c *Controller) snapAllForDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID) { - c.logger.Debugf("TURTLE: snapAllForDatabase: %s", qdbid) +func (c *Controller) snapAllForDatabase(tx dax.Transaction, qdbid dax.QualifiedDatabaseID, log logger.Logger) { + log.Debugf("snapAllForDatabase: %s", qdbid) computeNodes, err := c.Balancer.CurrentState(tx, dax.RoleTypeCompute, qdbid) if err != nil { - c.logger.Printf("Error getting compute balancer state for snapping turtle: %v", err) + log.Printf("Error getting compute balancer state for snapping turtle: %v", err) } // Weird nested loop for snapshotting shard data. The reason for @@ -72,10 +73,10 @@ func (c *Controller) snapAllForDatabase(tx dax.Transaction, qdbid dax.QualifiedD stillWorking = true j, err := decodeShard(workerInfo.Jobs[i]) if err != nil { - c.logger.Printf("couldn't decode a shard out of the job: '%s', err: %v", workerInfo.Jobs[i], err) + log.Printf("couldn't decode a shard out of the job: '%s', err: %v", workerInfo.Jobs[i], err) } if err := c.snapshotShardData(tx, j.t.QualifiedTableID(), j.shardNum()); err != nil { - c.logger.Printf("Couldn't snapshot table: %s, shard: %d, error: %v", j.t, j.shardNum(), err) + log.Printf("Couldn't snapshot table: %s, shard: %d, error: %v", j.t, j.shardNum(), err) } } i++ @@ -86,7 +87,7 @@ func (c *Controller) snapAllForDatabase(tx dax.Transaction, qdbid dax.QualifiedD // partitions. tables, err := c.Schemar.Tables(tx, dax.QualifiedDatabaseID{}) if err != nil { - c.logger.Printf("Couldn't get schema for snapshotting keys: %v", err) + log.Printf("Couldn't get schema for snapshotting keys: %v", err) return } // snapshot keyed fields @@ -96,7 +97,7 @@ func (c *Controller) snapAllForDatabase(tx dax.Transaction, qdbid dax.QualifiedD for _, f := range table.Fields { if f.StringKeys() && !f.IsPrimaryKey() { if err := c.snapshotFieldKeys(tx, table.QualifiedID(), f.Name); err != nil { - c.logger.Printf("Couldn't snapshot table: %s, field: %s, error: %v", table, f.Name, err) + log.Printf("Couldn't snapshot table: %s, field: %s, error: %v", table, f.Name, err) } } } @@ -108,7 +109,7 @@ func (c *Controller) snapAllForDatabase(tx dax.Transaction, qdbid dax.QualifiedD // back to back. translateNodes, err := c.Balancer.CurrentState(tx, dax.RoleTypeTranslate, qdbid) if err != nil { - c.logger.Printf("Error getting translate balancer state for snapping turtle: %v", err) + log.Printf("Error getting translate balancer state for snapping turtle: %v", err) } i = 0 @@ -125,13 +126,13 @@ func (c *Controller) snapAllForDatabase(tx dax.Transaction, qdbid dax.QualifiedD table := tableMap[j.table()] if table.StringKeys() { if err := c.snapshotTableKeys(tx, table.QualifiedID(), j.partitionNum()); err != nil { - c.logger.Printf("Couldn't snapshot table: %s, partition: %d, error: %v", table, j.partitionNum(), err) + log.Printf("Couldn't snapshot table: %s, partition: %d, error: %v", table, j.partitionNum(), err) } } - c.logger.Printf("couldn't decode a partition out of the job: '%s', err: %v", workerInfo.Jobs[i], err) + log.Printf("couldn't decode a partition out of the job: '%s', err: %v", workerInfo.Jobs[i], err) } } i++ } - c.logger.Debugf("TURTLE: snapAllForDatabase complete: %s", qdbid) + log.Debugf("snapAllForDatabase complete: %s", qdbid) } diff --git a/dax/queryer/service/queryer.go b/dax/queryer/service/queryer.go index 1d7d1f7b9..31eccc3cd 100644 --- a/dax/queryer/service/queryer.go +++ b/dax/queryer/service/queryer.go @@ -25,7 +25,7 @@ func New(uri *fbnet.URI, queryer *queryer.Queryer, logger logger.Logger) *querye return &queryerService{ uri: uri, queryer: queryer, - logger: logger, + logger: logger.WithPrefix("Queryer: "), } } diff --git a/dax/server/config.go b/dax/server/config.go index c41e7755b..a909102b3 100644 --- a/dax/server/config.go +++ b/dax/server/config.go @@ -71,7 +71,7 @@ func NewConfig() *Config { Config: controller.Config{ RegistrationBatchTimeout: time.Second * 3, StorageMethod: defaultStorageMethod, - SnappingTurtleTimeout: time.Second * 10, + SnappingTurtleTimeout: time.Minute * 3, }, }, Bind: ":" + defaultBindPort, diff --git a/dax/server/server.go b/dax/server/server.go index c682da3f7..eda997f8d 100644 --- a/dax/server/server.go +++ b/dax/server/server.go @@ -12,7 +12,6 @@ import ( "crypto/tls" "encoding/json" "io" - "log" "math/rand" "net" "os" @@ -73,7 +72,7 @@ func OptCommandConfig(config *Config) CommandOption { // OptCommandServiceManager allows the ability to pass in a ServiceManage that // has been initialized outside of the Command. This is useful for testing where -// we want to controll the service manager during a test run. +// we want to control the service manager during a test run. func OptCommandServiceManager(svcmgr *dax.ServiceManager) CommandOption { return func(c *Command) error { c.svcmgr = svcmgr @@ -111,7 +110,6 @@ func (m *Command) Start() (err error) { if seed == 0 { seed = time.Now().UTC().UnixNano() } - log.Printf("Random seed: %d", seed) rand.Seed(seed) if err := m.setupServer(); err != nil { @@ -217,7 +215,7 @@ func (m *Command) setupServer() error { if err != nil { return errors.Wrap(err, "marshalling config") } - m.logger.Printf("Config: %s", conf) + m.logger.Debugf("Config: %s", conf) // validateAddrs sets the appropriate values for Bind and Advertise // based on the inputs. It is not responsible for applying defaults, although diff --git a/dax/service_manager.go b/dax/service_manager.go index 04c0c7cee..9443c2fb1 100644 --- a/dax/service_manager.go +++ b/dax/service_manager.go @@ -100,8 +100,10 @@ func (s *ServiceManager) StopAll() error { // ControllerStart starts the Controller service. func (s *ServiceManager) ControllerStart() error { if s.Controller == nil { + s.Logger.Debugf("Skipping Controller") return nil } + s.Logger.Printf("Starting Controller") s.mu.Lock() defer s.mu.Unlock() @@ -126,6 +128,7 @@ func (s *ServiceManager) ControllerStop() error { if s.Controller == nil { return nil } + s.Logger.Printf("Stopping Controller") s.mu.Lock() defer s.mu.Unlock() @@ -148,8 +151,10 @@ func (s *ServiceManager) ControllerStop() error { // QueryerStart starts the Queryer service. func (s *ServiceManager) QueryerStart() error { if s.Queryer == nil { + s.Logger.Debugf("Skipping Queryer") return nil } + s.Logger.Printf("Starting Queryer") s.mu.Lock() defer s.mu.Unlock() @@ -174,6 +179,7 @@ func (s *ServiceManager) QueryerStop() error { if s.Queryer == nil { return nil } + s.Logger.Printf("Stopping Queryer") s.mu.Lock() defer s.mu.Unlock() @@ -207,6 +213,8 @@ func (s *ServiceManager) ComputerStart(key ServiceKey) error { s.mu.Lock() defer s.mu.Unlock() + s.Logger.Printf("Starting Computer: %s", key) + serviceState, ok := s.computers[key] if !ok { return errors.Errorf("computer to be started does not exist: %s", key) @@ -238,6 +246,8 @@ func (s *ServiceManager) ComputerStop(key ServiceKey) error { s.mu.Lock() defer s.mu.Unlock() + s.Logger.Printf("Stopping Computer: %s", key) + serviceState, ok := s.computers[key] if !ok { return errors.Errorf("computer to be stopped does not exist: %s", key) diff --git a/diagnostics_internal_test.go b/diagnostics_internal_test.go index 8cf9f336e..fb2490634 100644 --- a/diagnostics_internal_test.go +++ b/diagnostics_internal_test.go @@ -115,7 +115,7 @@ func TestDiagnosticsVersion_Check(t *testing.T) { // Create a new client. d := newDiagnosticsCollector("localhost:10101") - logs := logger.NewCaptureLogger() + logs := logger.NewBufferLogger() d.Logger = logs version := "0.1.1" @@ -126,11 +126,14 @@ func TestDiagnosticsVersion_Check(t *testing.T) { if err != nil { t.Fatalf("checking version: %v", err) } - if len(logs.Prints) != 1 { - t.Fatalf("expected a version upgrade message") + + allLogs, err := logs.ReadAll() + if err != nil { + t.Fatalf("reading all logs: %v", err) } - if !strings.Contains(logs.Prints[0], "a newer version") { - t.Fatalf("expected version upgrade message, got '%s'", logs.Prints[0]) + + if !strings.Contains(string(allLogs), "a newer version") { + t.Fatalf("expected version upgrade message, got '%s'", allLogs) } } diff --git a/idk/header_test.go b/idk/header_test.go index 02d4e3ed9..053b66f57 100644 --- a/idk/header_test.go +++ b/idk/header_test.go @@ -6,6 +6,8 @@ import ( "strings" "testing" "time" + + "github.com/featurebasedb/featurebase/v3/logger" ) func TestHeaderToField(t *testing.T) { @@ -421,3 +423,9 @@ func (r *recordingLogger) Printf(format string, v ...interface{}) { func (r *recordingLogger) Debugf(format string, v ...interface{}) { r.debugs = append(r.debugs, fmt.Sprintf(format, v...)) } + +// WithPrefix does nothing for the recording logger... just makes it +// implement the interface. Do not use. +func (r *recordingLogger) WithPrefix(prefix string) logger.Logger { + panic("unimplemented") +} diff --git a/idk/kafka/source_test.go b/idk/kafka/source_test.go index 27911b3a5..0738bc753 100644 --- a/idk/kafka/source_test.go +++ b/idk/kafka/source_test.go @@ -826,6 +826,7 @@ func TestRegistryURLParsing(t *testing.T) { } for i, test := range tests { + test := test t.Run(fmt.Sprintf("%s-%d", test.name, i), func(t *testing.T) { t.Parallel() diff --git a/idk/kinesis/logger.go b/idk/kinesis/logger.go index 08a36ead4..80898ea22 100644 --- a/idk/kinesis/logger.go +++ b/idk/kinesis/logger.go @@ -232,3 +232,7 @@ func (esl *ErrorStreamLogger) Panicf(format string, v ...interface{}) { esl.base.Errorf(errors.Wrap(err, errMsg).Error()) } } + +func (esl *ErrorStreamLogger) WithPrefix(prefix string) logger.Logger { + return NewErrorStreamLogger(esl.base.WithPrefix(prefix), esl.store) +} diff --git a/idk/kinesis/logger_test.go b/idk/kinesis/logger_test.go index 43586451e..14725d5ad 100644 --- a/idk/kinesis/logger_test.go +++ b/idk/kinesis/logger_test.go @@ -59,6 +59,10 @@ func (msl *MapStashLogger) Panicf(format string, v ...interface{}) { msl.messages["PANIC"] = append(msl.messages["PANIC"], fmt.Sprintf(format, v...)) } +func (msl *MapStashLogger) WithPrefix(prefix string) logger.Logger { + panic("unimplemented") +} + // In-memory error store implementation to use in unit tests. type MapStashErrorStore struct { storage map[ErrorType][]string diff --git a/logger/logger.go b/logger/logger.go index 60b0fce3c..204f5090a 100644 --- a/logger/logger.go +++ b/logger/logger.go @@ -27,6 +27,9 @@ type Logger interface { Warnf(format string, v ...interface{}) Errorf(format string, v ...interface{}) Panicf(format string, v ...interface{}) + // WithPrefix returns a new Logger with the same configuration as + // this one, but all logs will have the given prefix. + WithPrefix(prefix string) Logger } const ( @@ -66,10 +69,16 @@ func (n *nopLogger) Errorf(format string, v ...interface{}) {} // Panicf is a no-op implementation of the Logger Panicf method. func (n *nopLogger) Panicf(format string, v ...interface{}) {} +func (n *nopLogger) WithPrefix(prefix string) Logger { + return n +} + // standardLogger is a basic implementation of Logger based on log.Logger. type standardLogger struct { logger *log.Logger verbosity int + prefix string + w io.Writer } // write in UTC with constant width and microsecond resolution. @@ -81,21 +90,23 @@ func (fl formatLog) Write(bytes []byte) (int, error) { return fmt.Fprintf(fl.w, "%v %v", time.Now().UTC().Format(RFC3339UsecTz0), string(bytes)) } -func newStandardLogger(w io.Writer, verbosity int) *standardLogger { - logger := log.New(w, "", 0) +func newStandardLogger(w io.Writer, verbosity int, prefix string) *standardLogger { + logger := log.New(w, prefix, 0) logger.SetOutput(formatLog{w: w}) return &standardLogger{ logger: logger, verbosity: verbosity, + prefix: prefix, + w: w, } } func NewStandardLogger(w io.Writer) *standardLogger { - return newStandardLogger(w, LevelInfo) + return newStandardLogger(w, LevelInfo, "") } func NewVerboseLogger(w io.Writer) *standardLogger { - return newStandardLogger(w, LevelDebug) + return newStandardLogger(w, LevelDebug, "") } func (s *standardLogger) printf(level int, format string, v ...interface{}) { @@ -137,46 +148,8 @@ func (s *standardLogger) Logger() *log.Logger { return s.logger } -// CaptureLogger is a test logger that stores all the print and debug messages -// it sees. -type CaptureLogger struct { - Prints []string - Debugs []string -} - -// NewCaptureLogger yields a CaptureLogger. -func NewCaptureLogger() *CaptureLogger { - return &CaptureLogger{} -} - -// Printf formats a message and appends it to Debugs. -func (cl *CaptureLogger) Printf(format string, v ...interface{}) { - cl.Debugs = append(cl.Prints, fmt.Sprintf(LevelPrefix(LevelInfo)+format, v...)) -} - -// Debugf formats a message and appends it to Debugs. -func (cl *CaptureLogger) Debugf(format string, v ...interface{}) { - cl.Debugs = append(cl.Debugs, fmt.Sprintf(LevelPrefix(LevelDebug)+format, v...)) -} - -// Infof formats a message and appends it to Prints. -func (cl *CaptureLogger) Infof(format string, v ...interface{}) { - cl.Prints = append(cl.Prints, fmt.Sprintf(LevelPrefix(LevelInfo)+format, v...)) -} - -// Warnf formats a message and appends it to Prints. -func (cl *CaptureLogger) Warnf(format string, v ...interface{}) { - cl.Prints = append(cl.Prints, fmt.Sprintf(LevelPrefix(LevelWarn)+format, v...)) -} - -// Errorf formats a message and appends it to Prints. -func (cl *CaptureLogger) Errorf(format string, v ...interface{}) { - cl.Prints = append(cl.Prints, fmt.Sprintf(LevelPrefix(LevelError)+format, v...)) -} - -// Panicf formats a message and appends it to Prints. -func (cl *CaptureLogger) Panicf(format string, v ...interface{}) { - cl.Prints = append(cl.Prints, fmt.Sprintf(LevelPrefix(LevelPanic)+format, v...)) +func (s *standardLogger) WithPrefix(prefix string) Logger { + return newStandardLogger(s.w, s.verbosity, prefix) } // Logfer is a thing that has only a Logf() method, like for instance, @@ -215,6 +188,11 @@ func (ll *LogfLogger) Panicf(format string, v ...interface{}) { ll.wrapped.Logf(format, v...) } +// WithPrefix does nothing for LogfLogger because I'm lazy. +func (ll *LogfLogger) WithPrefix(prefix string) Logger { + return ll +} + func NewLogfLogger(l Logfer) *LogfLogger { return &LogfLogger{wrapped: l} } @@ -257,6 +235,11 @@ func (b *bufferLogger) Panicf(format string, v ...interface{}) { b.Printf(LevelPrefix(4)+format, v...) } +// WithPrefix does nothing for bufferLogger because I'm lazy. +func (b *bufferLogger) WithPrefix(prefix string) Logger { + return b +} + func (b *bufferLogger) ReadAll() ([]byte, error) { b.mu.Lock() defer b.mu.Unlock() diff --git a/server/config.go b/server/config.go index 8af0776a6..f2d42e8ac 100644 --- a/server/config.go +++ b/server/config.go @@ -382,7 +382,7 @@ func NewConfig() *Config { LongQueryTime: toml.Duration(-time.Minute), - CheckInInterval: 5 * time.Second, + CheckInInterval: 60 * time.Second, } // Cluster config. diff --git a/server/server.go b/server/server.go index d19c2249d..d9ecc9ea4 100644 --- a/server/server.go +++ b/server/server.go @@ -296,6 +296,7 @@ func (m *Command) StartNoServe(addr dax.Address) (err error) { // If the interval period is 0, the check-in is disabled. func (m *Command) checkIn(addr dax.Address) { interval := m.Config.CheckInInterval + m.logger.Printf("CheckInInterval: %v", interval) if interval == 0 { return diff --git a/server/server_test.go b/server/server_test.go index 0f2a6dfb8..88e60d803 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -44,6 +44,7 @@ func TestMain_Set_Quick(t *testing.T) { } for i := 0; i < 10; i++ { + i := i t.Run(fmt.Sprint(i), func(t *testing.T) { t.Parallel()