diff --git a/app/audit.go b/app/audit.go index deb1546bc5..e07282dd74 100644 --- a/app/audit.go +++ b/app/audit.go @@ -161,7 +161,9 @@ func (s *Server) configureAudit(adt *audit.Audit, bAllowAdvancedLogging bool) er errs = multierror.Append(err) continue } - adt.AddTarget(target) + if target != nil { + adt.AddTarget(target) + } } return errs } diff --git a/app/server.go b/app/server.go index 46c7f36344..6255817284 100644 --- a/app/server.go +++ b/app/server.go @@ -112,19 +112,20 @@ type Server struct { newStore func() store.Store - htmlTemplateWatcher *utils.HTMLTemplateWatcher - sessionCache cache.Cache - seenPendingPostIdsCache cache.Cache - statusCache cache.Cache - configListenerId string - licenseListenerId string - logListenerId string - clusterLeaderListenerId string - searchConfigListenerId string - searchLicenseListenerId string - configStore config.Store - asymmetricSigningKey *ecdsa.PrivateKey - postActionCookieSecret []byte + htmlTemplateWatcher *utils.HTMLTemplateWatcher + sessionCache cache.Cache + seenPendingPostIdsCache cache.Cache + statusCache cache.Cache + configListenerId string + licenseListenerId string + logListenerId string + clusterLeaderListenerId string + searchConfigListenerId string + searchLicenseListenerId string + loggerMetricsLicenseListenerId string + configStore config.Store + asymmetricSigningKey *ecdsa.PrivateKey + postActionCookieSecret []byte advancedLogListenerCleanup func() @@ -486,6 +487,11 @@ func NewServer(options ...Option) (*Server, error) { } } + s.enableLoggingMetrics() + s.loggerMetricsLicenseListenerId = s.AddLicenseListener(func(oldLicense, newLicense *model.License) { + s.enableLoggingMetrics() + }) + // Enable developer settings if this is a "dev" build if model.BuildNumber == "dev" { s.UpdateConfig(func(cfg *model.Config) { *cfg.ServiceSettings.EnableDeveloper = true }) @@ -617,6 +623,18 @@ func (s *Server) initLogging() error { return nil } +func (s *Server) enableLoggingMetrics() { + if s.Metrics == nil { + return + } + + if err := mlog.EnableMetrics(s.Metrics.GetLoggerMetricsCollector()); err != nil { + mlog.Debug("Failed to enable advanced logging metrics", mlog.Err(err)) + } else { + mlog.Debug("Advanced logging metrics enabled") + } +} + const TIME_TO_WAIT_FOR_CONNECTIONS_TO_CLOSE_ON_SERVER_SHUTDOWN = time.Second func (s *Server) StopHTTPServer() { @@ -649,6 +667,7 @@ func (s *Server) Shutdown() error { s.HubStop() s.ShutDownPlugins() s.RemoveLicenseListener(s.licenseListenerId) + s.RemoveLicenseListener(s.loggerMetricsLicenseListenerId) s.RemoveClusterLeaderChangedListener(s.clusterLeaderListenerId) if s.tracer != nil { diff --git a/einterfaces/metrics.go b/einterfaces/metrics.go index 04ad95c153..921fadca73 100644 --- a/einterfaces/metrics.go +++ b/einterfaces/metrics.go @@ -3,6 +3,8 @@ package einterfaces +import "github.com/mattermost/logr" + type MetricsInterface interface { StartServer() StopServer() @@ -56,4 +58,6 @@ type MetricsInterface interface { ObservePluginMultiHookIterationDuration(pluginID string, elapsed float64) ObservePluginMultiHookDuration(elapsed float64) ObservePluginApiDuration(pluginID, apiName string, success bool, elapsed float64) + + GetLoggerMetricsCollector() logr.MetricsCollector } diff --git a/einterfaces/mocks/MetricsInterface.go b/einterfaces/mocks/MetricsInterface.go index 746be9e242..b7b9564b06 100644 --- a/einterfaces/mocks/MetricsInterface.go +++ b/einterfaces/mocks/MetricsInterface.go @@ -4,7 +4,10 @@ package mocks -import mock "github.com/stretchr/testify/mock" +import ( + logr "github.com/mattermost/logr" + mock "github.com/stretchr/testify/mock" +) // MetricsInterface is an autogenerated mock type for the MetricsInterface type type MetricsInterface struct { @@ -31,6 +34,22 @@ func (_m *MetricsInterface) DecrementWebSocketBroadcastUsersRegistered(hub strin _m.Called(hub, amount) } +// GetLoggerMetricsCollector provides a mock function with given fields: +func (_m *MetricsInterface) GetLoggerMetricsCollector() logr.MetricsCollector { + ret := _m.Called() + + var r0 logr.MetricsCollector + if rf, ok := ret.Get(0).(func() logr.MetricsCollector); ok { + r0 = rf() + } else { + if ret.Get(0) != nil { + r0 = ret.Get(0).(logr.MetricsCollector) + } + } + + return r0 +} + // IncrementChannelIndexCounter provides a mock function with given fields: func (_m *MetricsInterface) IncrementChannelIndexCounter() { _m.Called() diff --git a/go.mod b/go.mod index cae17f0642..53e6f9b414 100644 --- a/go.mod +++ b/go.mod @@ -61,7 +61,7 @@ require ( github.com/mattermost/gorp v1.6.2-0.20200624165429-2595d5e54111 github.com/mattermost/gosaml2 v0.3.2 github.com/mattermost/ldap v0.0.0-20191128190019-9f62ba4b8d4d - github.com/mattermost/logr v1.0.5 + github.com/mattermost/logr v1.0.9 github.com/mattermost/rsc v0.0.0-20160330161541-bbaefb05eaa0 github.com/mattermost/viper v1.0.4 github.com/mattn/go-colorable v0.1.7 // indirect diff --git a/go.sum b/go.sum index ed4c3991be..2c153d4398 100644 --- a/go.sum +++ b/go.sum @@ -425,8 +425,10 @@ github.com/mattermost/gosaml2 v0.3.2 h1:kq2dY5qUe6fPPHra171GVlgo+ycBsEog0gZMetxL github.com/mattermost/gosaml2 v0.3.2/go.mod h1:Z429EIOiEi9kbq6yHoApfzlcXpa6dzRDc6pO+Vy2Ksk= github.com/mattermost/ldap v0.0.0-20191128190019-9f62ba4b8d4d h1:2DV7VIlEv6J5R5o6tUcb3ZMKJYeeZuWZL7Rv1m23TgQ= github.com/mattermost/ldap v0.0.0-20191128190019-9f62ba4b8d4d/go.mod h1:HLbgMEI5K131jpxGazJ97AxfPDt31osq36YS1oxFQPQ= -github.com/mattermost/logr v1.0.5 h1:TST38xROPguNh8o90BfDHpp1bz6HfTdFYX5Btw/oLwM= -github.com/mattermost/logr v1.0.5/go.mod h1:YzldchiJXgF789YNDFGXVoCHTQOTrCKwWft9Fwev1iI= +github.com/mattermost/logr v1.0.8 h1:Frkfo+FXGOShiSN3pljflqEMS9AKGySdTHHu36MYPXU= +github.com/mattermost/logr v1.0.8/go.mod h1:Mt4DPu1NXMe6JxPdwCC0XBoxXmN9eXOIRPoZarU2PXs= +github.com/mattermost/logr v1.0.9 h1:jw6f6CjPC2YfPqzpGVM/vCGuqKLJdVS400ZRTFtEQwQ= +github.com/mattermost/logr v1.0.9/go.mod h1:Mt4DPu1NXMe6JxPdwCC0XBoxXmN9eXOIRPoZarU2PXs= github.com/mattermost/rsc v0.0.0-20160330161541-bbaefb05eaa0 h1:G9tL6JXRBMzjuD1kkBtcnd42kUiT6QDwxfFYu7adM6o= github.com/mattermost/rsc v0.0.0-20160330161541-bbaefb05eaa0/go.mod h1:nV5bfVpT//+B1RPD2JvRnxbkLmJEYXmRaaVl15fsXjs= github.com/mattermost/viper v1.0.4 h1:cMYOz4PhguscGSPxrSokUtib5HrG4gCpiUh27wyA3d0= diff --git a/mlog/default.go b/mlog/default.go index 3b3829fdeb..6322650d6b 100644 --- a/mlog/default.go +++ b/mlog/default.go @@ -81,3 +81,9 @@ func defaultAddTarget(target logr.Target) error { // logger is replaced with mlog.Logger instance. return errors.New("cannot AddTarget on default logger") } + +func defaultEnableMetrics(collector logr.MetricsCollector) error { + // mlog.EnableMetrics should not be called until default + // logger is replaced with mlog.Logger instance. + return errors.New("cannot EnableMetrics on default logger") +} diff --git a/mlog/global.go b/mlog/global.go index 5eec280be9..39b4b30e2f 100644 --- a/mlog/global.go +++ b/mlog/global.go @@ -32,6 +32,7 @@ func InitGlobalLogger(logger *Logger) { ConfigAdvancedLogging = globalLogger.ConfigAdvancedLogging ShutdownAdvancedLogging = globalLogger.ShutdownAdvancedLogging AddTarget = globalLogger.AddTarget + EnableMetrics = globalLogger.EnableMetrics } func RedirectStdLog(logger *Logger) { @@ -45,6 +46,7 @@ type FlushFunc func(context.Context) error type ConfigFunc func(cfg LogTargetCfg) error type ShutdownFunc func(context.Context) error type AddTargetFunc func(logr.Target) error +type EnableMetricsFunc func(logr.MetricsCollector) error // DON'T USE THIS Modify the level on the app logger func GloballyDisableDebugLogForTest() { @@ -68,3 +70,4 @@ var Flush FlushFunc = defaultFlush var ConfigAdvancedLogging ConfigFunc = defaultAdvancedConfig var ShutdownAdvancedLogging ShutdownFunc = defaultAdvancedShutdown var AddTarget AddTargetFunc = defaultAddTarget +var EnableMetrics EnableMetricsFunc = defaultEnableMetrics diff --git a/mlog/log.go b/mlog/log.go index f2e99a128c..c314a5732f 100644 --- a/mlog/log.go +++ b/mlog/log.go @@ -87,6 +87,7 @@ func NewLogger(config *LoggerConfiguration) *Logger { logger := &Logger{ consoleLevel: zap.NewAtomicLevelAt(getZapLevel(config.ConsoleLevel)), fileLevel: zap.NewAtomicLevelAt(getZapLevel(config.FileLevel)), + logrLogger: newLogr(), } if config.EnableConsole { @@ -101,6 +102,7 @@ func NewLogger(config *LoggerConfiguration) *Logger { MaxSize: 100, Compress: true, }) + core := zapcore.NewCore(makeEncoder(config.FileJson), writer, logger.fileLevel) cores = append(cores, core) } @@ -167,77 +169,67 @@ func (l *Logger) Sugar() *SugarLogger { func (l *Logger) Debug(message string, fields ...Field) { l.zap.Debug(message, fields...) - if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Debug) { + if isLevelEnabled(l.logrLogger, logr.Debug) { l.logrLogger.WithFields(zapToLogr(fields)).Debug(message) } } func (l *Logger) Info(message string, fields ...Field) { l.zap.Info(message, fields...) - if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Info) { + if isLevelEnabled(l.logrLogger, logr.Info) { l.logrLogger.WithFields(zapToLogr(fields)).Info(message) } } func (l *Logger) Warn(message string, fields ...Field) { l.zap.Warn(message, fields...) - if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Warn) { + if isLevelEnabled(l.logrLogger, logr.Warn) { l.logrLogger.WithFields(zapToLogr(fields)).Warn(message) } } func (l *Logger) Error(message string, fields ...Field) { l.zap.Error(message, fields...) - if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Error) { + if isLevelEnabled(l.logrLogger, logr.Error) { l.logrLogger.WithFields(zapToLogr(fields)).Error(message) } } func (l *Logger) Critical(message string, fields ...Field) { l.zap.Error(message, fields...) - if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Error) { + if isLevelEnabled(l.logrLogger, logr.Error) { l.logrLogger.WithFields(zapToLogr(fields)).Error(message) } } func (l *Logger) Log(level LogLevel, message string, fields ...Field) { - if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Level(level)) { - l.logrLogger.WithFields(zapToLogr(fields)).Log(logr.Level(level), message) - } + l.logrLogger.WithFields(zapToLogr(fields)).Log(logr.Level(level), message) } func (l *Logger) LogM(levels []LogLevel, message string, fields ...Field) { - if l.logrLogger != nil { - var logger *logr.Logger - for _, lvl := range levels { - if isLevelEnabled(l.logrLogger, logr.Level(lvl)) { - // don't create logger with fields unless at least one level is active. - if logger == nil { - l := l.logrLogger.WithFields(zapToLogr(fields)) - logger = &l - } - logger.Log(logr.Level(lvl), message) + var logger *logr.Logger + for _, lvl := range levels { + if isLevelEnabled(l.logrLogger, logr.Level(lvl)) { + // don't create logger with fields unless at least one level is active. + if logger == nil { + l := l.logrLogger.WithFields(zapToLogr(fields)) + logger = &l } + logger.Log(logr.Level(lvl), message) } } } func (l *Logger) Flush(cxt context.Context) error { - if l.logrLogger != nil { - return l.logrLogger.Logr().Flush() // TODO: use context when Logr lib supports it. - } - return nil + return l.logrLogger.Logr().Flush() // TODO: use context when Logr lib supports it. } // ShutdownAdvancedLogging stops the logger from accepting new log records and tries to // flush queues within the context timeout. Once complete all targets are shutdown // and any resources released. func (l *Logger) ShutdownAdvancedLogging(cxt context.Context) error { - var err error - if l.logrLogger != nil { - err = l.logrLogger.Logr().Shutdown() // TODO: use context when Logr lib supports it. - l.logrLogger = nil - } + err := l.logrLogger.Logr().Shutdown() // TODO: use context when Logr lib supports it. + l.logrLogger = newLogr() return err } @@ -245,14 +237,11 @@ func (l *Logger) ShutdownAdvancedLogging(cxt context.Context) error { // specified log targets. This is the easiest way to get the advanced logger // configured via a config source such as file. func (l *Logger) ConfigAdvancedLogging(targets LogTargetCfg) error { - if l.logrLogger != nil { - if err := l.ShutdownAdvancedLogging(context.Background()); err != nil { - Error("error shutting down previous logger", Err(err)) - } + if err := l.ShutdownAdvancedLogging(context.Background()); err != nil { + Error("error shutting down previous logger", Err(err)) } - logr, err := newLogr(targets) - l.logrLogger = logr + err := logrAddTargets(l.logrLogger, targets) return err } @@ -262,3 +251,9 @@ func (l *Logger) ConfigAdvancedLogging(targets LogTargetCfg) error { func (l *Logger) AddTarget(target logr.Target) error { return l.logrLogger.Logr().AddTarget(target) } + +// EnableMetrics enables metrics collection by supplying a MetricsCollector. +// The MetricsCollector provides counters and gauges that are updated by log targets. +func (l *Logger) EnableMetrics(collector logr.MetricsCollector) error { + return l.logrLogger.Logr().SetMetricsCollector(collector) +} diff --git a/mlog/logr.go b/mlog/logr.go index 65248c49e1..dba08e1006 100644 --- a/mlog/logr.go +++ b/mlog/logr.go @@ -28,7 +28,7 @@ type LogLevel struct { } type LogTarget struct { - Type string // one of "console", "file", "tcp", "syslog". + Type string // one of "console", "file", "tcp", "syslog", "none". Format string // one of "json", "plain" Levels []LogLevel Options json.RawMessage @@ -38,7 +38,7 @@ type LogTarget struct { type LogTargetCfg map[string]*LogTarget type LogrCleanup func() error -func newLogr(targets LogTargetCfg) (*logr.Logger, error) { +func newLogr() *logr.Logger { lgr := &logr.Logr{} lgr.OnExit = func(int) {} lgr.OnPanic = func(interface{}) {} @@ -46,12 +46,12 @@ func newLogr(targets LogTargetCfg) (*logr.Logger, error) { lgr.OnQueueFull = onQueueFull lgr.OnTargetQueueFull = onTargetQueueFull - err := logrAddTargets(lgr, targets) logger := lgr.NewLogger() - return &logger, err + return &logger } -func logrAddTargets(lgr *logr.Logr, targets LogTargetCfg) error { +func logrAddTargets(logger *logr.Logger, targets LogTargetCfg) error { + lgr := logger.Logr() var errs error for name, t := range targets { target, err := NewLogrTarget(name, t) @@ -59,7 +59,9 @@ func logrAddTargets(lgr *logr.Logr, targets LogTargetCfg) error { errs = multierror.Append(err) continue } - lgr.AddTarget(target) + if target != nil { + lgr.AddTarget(target) + } } return errs } @@ -90,6 +92,8 @@ func NewLogrTarget(name string, t *LogTarget) (logr.Target, error) { return newSyslogTarget(name, t, filter, formatter) case "tcp": return newTCPTarget(name, t, filter, formatter) + case "none": + return nil, nil } return nil, fmt.Errorf("invalid type '%s' for target %s", t.Type, name) } @@ -201,6 +205,10 @@ func checkFileWritable(filename string) error { } func isLevelEnabled(logger *logr.Logger, level logr.Level) bool { + if logger == nil || logger.Logr() == nil { + return false + } + status := logger.Logr().IsLevelEnabled(level) return status.Enabled } diff --git a/mlog/tcp_test.go b/mlog/tcp_test.go index 111dc085e2..0eb689385b 100644 --- a/mlog/tcp_test.go +++ b/mlog/tcp_test.go @@ -38,15 +38,16 @@ func TestNewTcpTarget(t *testing.T) { data := []string{"I drink your milkshake!", "We don't need no badges!", "You can't fight in here! This is the war room!"} - logr, err := newLogr(targets) + logger := newLogr() + err = logrAddTargets(logger, targets) require.NoError(t, err) for _, s := range data { - logr.Info(s) + logger.Info(s) } - err = logr.Logr().Flush() + err = logger.Logr().Flush() require.NoError(t, err) - err = logr.Logr().Shutdown() + err = logger.Logr().Shutdown() require.NoError(t, err) err = server.waitForAnyConnection() diff --git a/mlog/testing.go b/mlog/testing.go index bf1bcedf43..f66c0f4d45 100644 --- a/mlog/testing.go +++ b/mlog/testing.go @@ -32,6 +32,7 @@ func NewTestingLogger(tb testing.TB, writer io.Writer) *Logger { testingLogger := &Logger{ consoleLevel: zap.NewAtomicLevelAt(getZapLevel("debug")), fileLevel: zap.NewAtomicLevelAt(getZapLevel("info")), + logrLogger: newLogr(), } logWriterCore := zapcore.NewCore(makeEncoder(true), logWriterSync, testingLogger.consoleLevel) diff --git a/vendor/github.com/mattermost/logr/go.mod b/vendor/github.com/mattermost/logr/go.mod index 0feeb6eb08..e8e8acfb2f 100644 --- a/vendor/github.com/mattermost/logr/go.mod +++ b/vendor/github.com/mattermost/logr/go.mod @@ -4,6 +4,7 @@ go 1.12 require ( github.com/francoispqt/gojay v1.2.13 + github.com/stretchr/testify v1.2.2 github.com/wiggin77/cfg v1.0.2 github.com/wiggin77/merror v1.0.2 gopkg.in/natefinch/lumberjack.v2 v2.0.0 diff --git a/vendor/github.com/mattermost/logr/logr.go b/vendor/github.com/mattermost/logr/logr.go index a293c16bc7..9cda6d1cf0 100644 --- a/vendor/github.com/mattermost/logr/logr.go +++ b/vendor/github.com/mattermost/logr/logr.go @@ -27,6 +27,13 @@ type Logr struct { shutdown bool lvlCache levelCache + metricsOnce sync.Once + metricsDone chan struct{} + metrics MetricsCollector + queueSizeGauge Gauge + loggedCounter Counter + errorCounter Counter + bufferPool sync.Pool // MaxQueueSize is the maximum number of log records that can be queued. @@ -95,6 +102,10 @@ type Logr struct { // DisableBufferPool when true disables the buffer pool. See MaxPooledBuffer. DisableBufferPool bool + + // MetricsUpdateFreqMillis determines how often polled metrics are updated + // when metrics are enabled. + MetricsUpdateFreqMillis int64 } // Configure adds/removes targets via the supplied `Config`. @@ -117,6 +128,13 @@ func (logr *Logr) AddTarget(target Target) error { defer logr.tmux.Unlock() logr.targets = append(logr.targets, target) + var err error + if logr.metrics != nil { + if tm, ok := target.(TargetWithMetrics); ok { + err = tm.EnableMetrics(logr.metrics, logr.MetricsUpdateFreqMillis) + } + } + logr.once.Do(func() { logr.maxQueueSizeActual = logr.MaxQueueSize if logr.maxQueueSizeActual == 0 { @@ -144,7 +162,7 @@ func (logr *Logr) AddTarget(target Target) error { go logr.start() }) logr.resetLevelCache() - return nil + return err } // NewLogger creates a Logger using defaults. A `Logger` is light-weight @@ -201,6 +219,13 @@ func (logr *Logr) IsLevelEnabled(lvl Level) LevelStatus { return status } +// HasTargets returns true only if at least one target exists within the Logr. +func (logr *Logr) HasTargets() bool { + logr.tmux.RLock() + defer logr.tmux.RUnlock() + return len(logr.targets) > 0 +} + // ResetLevelCache resets the cached results of `IsLevelEnabled`. This is // called any time a Target is added or a target's level is changed. func (logr *Logr) ResetLevelCache() { @@ -279,6 +304,10 @@ func (logr *Logr) panic(err interface{}) { // timing out. Use `IsTimeoutError` to determine if the returned error is // due to a timeout. func (logr *Logr) Flush() error { + if !logr.HasTargets() { + return nil + } + logr.mux.Lock() defer logr.mux.Unlock() @@ -310,6 +339,10 @@ func (logr *Logr) Shutdown() error { } logr.shutdown = true logr.resetLevelCache() + if logr.metricsDone != nil { + close(logr.metricsDone) + logr.metricsDone = nil + } logr.mux.Unlock() errs := merror.New() @@ -344,6 +377,9 @@ func (logr *Logr) Shutdown() error { // If `OnLoggerError` is not nil, it is called with the error, otherwise the error is // output to `os.Stderr`. func (logr *Logr) ReportError(err interface{}) { + if logr.errorCounter != nil { + logr.errorCounter.Inc() + } if logr.OnLoggerError == nil { fmt.Fprintln(os.Stderr, err) return @@ -415,6 +451,29 @@ func (logr *Logr) start() { close(logr.done) } +// startMetricsUpdater updates the metrics for any polled values every `MetricsUpdateFreqSecs` seconds until +// logr is closed. +func (logr *Logr) startMetricsUpdater() { + for { + updateFreq := logr.MetricsUpdateFreqMillis + if updateFreq == 0 { + updateFreq = DefMetricsUpdateFreqMillis + } + if updateFreq < 250 { + updateFreq = 250 // don't peg the CPU + } + + select { + case <-logr.metricsDone: + return + case <-time.After(time.Duration(updateFreq) * time.Millisecond): + if logr.queueSizeGauge != nil { + logr.queueSizeGauge.Set(float64(len(logr.in))) + } + } + } +} + // fanout pushes a LogRec to all targets. func (logr *Logr) fanout(rec *LogRec) { var target Target @@ -424,13 +483,20 @@ func (logr *Logr) fanout(rec *LogRec) { } }() + var logged bool + logr.tmux.RLock() defer logr.tmux.RUnlock() for _, target = range logr.targets { if enabled, _ := target.IsLevelEnabled(rec.Level()); enabled { target.Log(rec) + logged = true } } + + if logged && logr.loggedCounter != nil { + logr.loggedCounter.Inc() + } } // flush drains the queue and notifies when done. diff --git a/vendor/github.com/mattermost/logr/metrics.go b/vendor/github.com/mattermost/logr/metrics.go new file mode 100644 index 0000000000..ac992b243b --- /dev/null +++ b/vendor/github.com/mattermost/logr/metrics.go @@ -0,0 +1,85 @@ +package logr + +import ( + "errors" + + "github.com/wiggin77/merror" +) + +const ( + DefMetricsUpdateFreqMillis = 15000 // 15 seconds +) + +// Counter is a simple metrics sink that can only increment a value. +// Implementations are external to Logr and provided via `MetricsCollector`. +type Counter interface { + // Inc increments the counter by 1. Use Add to increment it by arbitrary non-negative values. + Inc() + // Add adds the given value to the counter. It panics if the value is < 0. + Add(float64) +} + +// Gauge is a simple metrics sink that can receive values and increase or decrease. +// Implementations are external to Logr and provided via `MetricsCollector`. +type Gauge interface { + // Set sets the Gauge to an arbitrary value. + Set(float64) + // Add adds the given value to the Gauge. (The value can be negative, resulting in a decrease of the Gauge.) + Add(float64) + // Sub subtracts the given value from the Gauge. (The value can be negative, resulting in an increase of the Gauge.) + Sub(float64) +} + +// MetricsCollector provides a way for users of this Logr package to have metrics pushed +// in an efficient way to any backend, e.g. Prometheus. +// For each target added to Logr, the supplied MetricsCollector will provide a Gauge +// and Counters that will be called frequently as logging occurs. +type MetricsCollector interface { + // QueueSizeGauge returns a Gauge that will be updated by the named target. + QueueSizeGauge(target string) (Gauge, error) + // LoggedCounter returns a Counter that will be incremented by the named target. + LoggedCounter(target string) (Counter, error) + // ErrorCounter returns a Counter that will be incremented by the named target. + ErrorCounter(target string) (Counter, error) + // DroppedCounter returns a Counter that will be incremented by the named target. + DroppedCounter(target string) (Counter, error) + // BlockedCounter returns a Counter that will be incremented by the named target. + BlockedCounter(target string) (Counter, error) +} + +// TargetWithMetrics is a target that provides metrics. +type TargetWithMetrics interface { + EnableMetrics(collector MetricsCollector, updateFreqMillis int64) error +} + +// SetMetricsCollector enables metrics collection by supplying a MetricsCollector. +// The MetricsCollector provides counters and gauges that are updated by log targets. +func (logr *Logr) SetMetricsCollector(collector MetricsCollector) error { + if collector == nil { + return errors.New("collector cannot be nil") + } + + logr.metrics = collector + logr.queueSizeGauge, _ = collector.QueueSizeGauge("_logr") + logr.loggedCounter, _ = collector.LoggedCounter("_logr") + logr.errorCounter, _ = collector.ErrorCounter("_logr") + + logr.metricsOnce.Do(func() { + logr.metricsDone = make(chan struct{}) + go logr.startMetricsUpdater() + }) + + merr := merror.New() + + logr.tmux.RLock() + defer logr.tmux.RUnlock() + for _, target := range logr.targets { + if tm, ok := target.(TargetWithMetrics); ok { + if err := tm.EnableMetrics(logr.metrics, logr.MetricsUpdateFreqMillis); err != nil { + merr.Append(err) + } + } + + } + return merr.ErrorOrNil() +} diff --git a/vendor/github.com/mattermost/logr/target.go b/vendor/github.com/mattermost/logr/target.go index bab71ec209..2ce76333cf 100644 --- a/vendor/github.com/mattermost/logr/target.go +++ b/vendor/github.com/mattermost/logr/target.go @@ -10,6 +10,9 @@ import ( // Target represents a destination for log records such as file, // database, TCP socket, etc. type Target interface { + // SetName provides an option name for the target. + SetName(name string) + // IsLevelEnabled returns true if this target should emit // logs for the specified level. Also determines if // a stack trace is required. @@ -33,9 +36,10 @@ type RecordWriter interface { // Basic provides the basic functionality of a Target that can be used // to more easily compose your own Targets. To use, just embed Basic -// in your target type, implement `RecordWriter`, and call `Start`. +// in your target type, implement `RecordWriter`, and call `(*Basic).Start`. type Basic struct { target Target + name string filter Filter formatter Formatter @@ -43,6 +47,14 @@ type Basic struct { in chan *LogRec done chan struct{} w RecordWriter + + queueSizeGauge Gauge + loggedCounter Counter + errorCounter Counter + droppedCounter Counter + blockedCounter Counter + + metricsUpdateFreqMillis int64 } // Start initializes this target helper and starts accepting log records for processing. @@ -61,6 +73,14 @@ func (b *Basic) Start(target Target, rw RecordWriter, filter Filter, formatter F b.done = make(chan struct{}, 1) b.w = rw go b.start() + + if b.queueSizeGauge != nil { + go b.startMetricsUpdater() + } +} + +func (b *Basic) SetName(name string) { + b.name = name } // IsLevelEnabled returns true if this target should emit @@ -97,8 +117,15 @@ func (b *Basic) Log(rec *LogRec) { default: handler := lgr.OnTargetQueueFull if handler != nil && handler(b.target, rec, cap(b.in)) { + if b.droppedCounter != nil { + b.droppedCounter.Inc() + } return // drop the record } + if b.blockedCounter != nil { + b.blockedCounter.Inc() + } + select { case <-time.After(lgr.enqueueTimeout()): lgr.ReportError(fmt.Errorf("target enqueue timeout for log rec [%v]", rec)) @@ -107,6 +134,39 @@ func (b *Basic) Log(rec *LogRec) { } } +// Metrics enables metrics collection using the provided MetricsCollector. +func (b *Basic) EnableMetrics(collector MetricsCollector, updateFreqMillis int64) error { + b.metricsUpdateFreqMillis = updateFreqMillis + + name := fmt.Sprintf("%v", b) + var err error + + if b.queueSizeGauge, err = collector.QueueSizeGauge(name); err != nil { + return err + } + if b.loggedCounter, err = collector.LoggedCounter(name); err != nil { + return err + } + if b.errorCounter, err = collector.ErrorCounter(name); err != nil { + return err + } + if b.droppedCounter, err = collector.DroppedCounter(name); err != nil { + return err + } + if b.blockedCounter, err = collector.BlockedCounter(name); err != nil { + return err + } + return nil +} + +// String returns a name for this target. Use `SetName` to specify a name. +func (b *Basic) String() string { + if b.name != "" { + return b.name + } + return fmt.Sprintf("%T", b.target) +} + // Start accepts log records via In channel and writes to the // supplied writer, until Done channel signaled. func (b *Basic) start() { @@ -123,13 +183,41 @@ func (b *Basic) start() { } else { err := b.w.Write(rec) if err != nil { + if b.errorCounter != nil { + b.errorCounter.Inc() + } rec.Logger().Logr().ReportError(err) + } else if b.loggedCounter != nil { + b.loggedCounter.Inc() } } } close(b.done) } +// startMetricsUpdater updates the metrics for any polled values every `MetricsUpdateFreqSecs` seconds until +// target is closed. +func (b *Basic) startMetricsUpdater() { + for { + updateFreq := b.metricsUpdateFreqMillis + if updateFreq == 0 { + updateFreq = DefMetricsUpdateFreqMillis + } + if updateFreq < 250 { + updateFreq = 250 // don't peg the CPU + } + + select { + case <-b.done: + return + case <-time.After(time.Duration(updateFreq) * time.Millisecond): + if b.queueSizeGauge != nil { + b.queueSizeGauge.Set(float64(len(b.in))) + } + } + } +} + // flush drains the queue and notifies when done. func (b *Basic) flush(done chan<- struct{}) { for { @@ -141,6 +229,9 @@ func (b *Basic) flush(done chan<- struct{}) { if rec.flush == nil { err = b.w.Write(rec) if err != nil { + if b.errorCounter != nil { + b.errorCounter.Inc() + } rec.Logger().Logr().ReportError(err) } } diff --git a/vendor/github.com/mattermost/logr/target/file.go b/vendor/github.com/mattermost/logr/target/file.go index bc0bcd1707..0fd50768da 100644 --- a/vendor/github.com/mattermost/logr/target/file.go +++ b/vendor/github.com/mattermost/logr/target/file.go @@ -85,8 +85,3 @@ func (f *File) Shutdown(ctx context.Context) error { return errs.ErrorOrNil() } - -// String returns a string representation of this target. -func (f *File) String() string { - return "FileTarget" -} diff --git a/vendor/github.com/mattermost/logr/target/syslog.go b/vendor/github.com/mattermost/logr/target/syslog.go index 2258fd29f5..1d2013b681 100644 --- a/vendor/github.com/mattermost/logr/target/syslog.go +++ b/vendor/github.com/mattermost/logr/target/syslog.go @@ -87,8 +87,3 @@ func (s *Syslog) Write(rec *logr.LogRec) error { } return err } - -// String returns a string representation of this target. -func (s *Syslog) String() string { - return "SyslogTarget" -} diff --git a/vendor/github.com/mattermost/logr/target/writer.go b/vendor/github.com/mattermost/logr/target/writer.go index b12b476046..2250da5138 100644 --- a/vendor/github.com/mattermost/logr/target/writer.go +++ b/vendor/github.com/mattermost/logr/target/writer.go @@ -38,8 +38,3 @@ func (w *Writer) Write(rec *logr.LogRec) error { _, err = w.out.Write(buf.Bytes()) return err } - -// String returns a string representation of this target. -func (w *Writer) String() string { - return "WriterTarget" -} diff --git a/vendor/modules.txt b/vendor/modules.txt index ad23d814af..e5cc06ff26 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -282,7 +282,7 @@ github.com/mattermost/gosaml2/uuid # github.com/mattermost/ldap v0.0.0-20191128190019-9f62ba4b8d4d ## explicit github.com/mattermost/ldap -# github.com/mattermost/logr v1.0.5 +# github.com/mattermost/logr v1.0.9 ## explicit github.com/mattermost/logr github.com/mattermost/logr/format