diff --git a/api4/main_test.go b/api4/main_test.go index e136f23931..403c22b5db 100644 --- a/api4/main_test.go +++ b/api4/main_test.go @@ -6,6 +6,7 @@ package api4 import ( "testing" + "github.com/mattermost/mattermost-server/v5/mlog" "github.com/mattermost/mattermost-server/v5/testlib" ) @@ -15,6 +16,8 @@ func TestMain(m *testing.M) { EnableResources: true, } + mlog.DisableZap() + mainHelper = testlib.NewMainHelperWithOptions(&options) defer mainHelper.Close() diff --git a/app/admin.go b/app/admin.go index aad393b027..245753f72e 100644 --- a/app/admin.go +++ b/app/admin.go @@ -4,6 +4,7 @@ package app import ( + "context" "fmt" "io" "io/ioutil" @@ -59,6 +60,10 @@ func (s *Server) GetLogsSkipSend(page, perPage int) ([]string, *model.AppError) var lines []string if *s.Config().LogSettings.EnableFile { + timeoutCtx, timeoutCancel := context.WithTimeout(context.Background(), mlog.DefaultFlushTimeout) + defer timeoutCancel() + mlog.Flush(timeoutCtx) + logFile := utils.GetLogFileLocation(*s.Config().LogSettings.FileLocation) file, err := os.Open(logFile) if err != nil { diff --git a/app/main_test.go b/app/main_test.go index 562b676422..1a24e6d267 100644 --- a/app/main_test.go +++ b/app/main_test.go @@ -6,6 +6,7 @@ package app import ( "testing" + "github.com/mattermost/mattermost-server/v5/mlog" "github.com/mattermost/mattermost-server/v5/testlib" ) @@ -17,6 +18,8 @@ func TestMain(m *testing.M) { EnableResources: true, } + mlog.DisableZap() + mainHelper = testlib.NewMainHelperWithOptions(&options) defer mainHelper.Close() diff --git a/app/server.go b/app/server.go index c829dc5885..bbaab11457 100644 --- a/app/server.go +++ b/app/server.go @@ -115,20 +115,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 - loggerMetricsLicenseListenerId 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 + loggerLicenseListenerId string + configStore config.Store + asymmetricSigningKey *ecdsa.PrivateKey + postActionCookieSecret []byte advancedLogListenerCleanup func() @@ -480,19 +480,11 @@ func NewServer(options ...Option) (*Server, error) { } } - if !allowAdvancedLogging { - timeoutCtx, cancelCtx := context.WithTimeout(context.Background(), time.Second*5) - defer cancelCtx() - mlog.Info("Shutting down advanced logging") - mlog.ShutdownAdvancedLogging(timeoutCtx) - if s.advancedLogListenerCleanup != nil { - s.advancedLogListenerCleanup() - s.advancedLogListenerCleanup = nil - } - } - + s.removeUnlicensedLogTargets(license) s.enableLoggingMetrics() - s.loggerMetricsLicenseListenerId = s.AddLicenseListener(func(oldLicense, newLicense *model.License) { + + s.loggerLicenseListenerId = s.AddLicenseListener(func(oldLicense, newLicense *model.License) { + s.removeUnlicensedLogTargets(newLicense) s.enableLoggingMetrics() }) @@ -561,6 +553,7 @@ func (s *Server) AppOptions() []AppOption { } } +// initLogging initializes and configures the logger. This may be called more than once. func (s *Server) initLogging() error { if s.Log == nil { s.Log = mlog.NewLogger(utils.MloggerConfigFromLoggerConfig(&s.Config().LogSettings, utils.GetLogFileLocation)) @@ -578,6 +571,9 @@ func (s *Server) initLogging() error { // Use this app logger as the global logger (eventually remove all instances of global logging) mlog.InitGlobalLogger(s.Log) + if s.logListenerId != "" { + s.RemoveConfigListener(s.logListenerId) + } s.logListenerId = s.AddConfigListener(func(_, after *model.Config) { s.Log.ChangeLevels(utils.MloggerConfigFromLoggerConfig(&after.LogSettings, utils.GetLogFileLocation)) @@ -633,13 +629,27 @@ func (s *Server) initLogging() error { return nil } +func (s *Server) removeUnlicensedLogTargets(license *model.License) { + if license != nil && *license.Features.AdvancedLogging { + // advanced logging enabled via license; no need to remove any targets + return + } + + timeoutCtx, cancelCtx := context.WithTimeout(context.Background(), time.Second*10) + defer cancelCtx() + + mlog.RemoveTargets(timeoutCtx, func(ti mlog.TargetInfo) bool { + return ti.Type != "*target.Writer" && ti.Type != "*target.File" + }) +} + 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)) + mlog.Error("Failed to enable advanced logging metrics", mlog.Err(err)) } else { mlog.Debug("Advanced logging metrics enabled") } @@ -677,7 +687,7 @@ func (s *Server) Shutdown() error { s.HubStop() s.ShutDownPlugins() s.RemoveLicenseListener(s.licenseListenerId) - s.RemoveLicenseListener(s.loggerMetricsLicenseListenerId) + s.RemoveLicenseListener(s.loggerLicenseListenerId) s.RemoveClusterLeaderChangedListener(s.clusterLeaderListenerId) if s.tracer != nil { diff --git a/app/server_test.go b/app/server_test.go index 56663f60b2..c97f9be781 100644 --- a/app/server_test.go +++ b/app/server_test.go @@ -304,6 +304,10 @@ func TestPanicLog(t *testing.T) { require.NoError(t, os.Remove(tmpfile.Name())) }() + // This test requires Zap file target for now. + mlog.EnableZap() + defer mlog.DisableZap() + // Creating logger to log to console and temp file logger := mlog.NewLogger(&mlog.LoggerConfiguration{ EnableConsole: true, diff --git a/app/slashcommands/main_test.go b/app/slashcommands/main_test.go index 220abc147c..03aa15ac40 100644 --- a/app/slashcommands/main_test.go +++ b/app/slashcommands/main_test.go @@ -6,6 +6,7 @@ package slashcommands import ( "testing" + "github.com/mattermost/mattermost-server/v5/mlog" "github.com/mattermost/mattermost-server/v5/testlib" ) @@ -17,6 +18,8 @@ func TestMain(m *testing.M) { EnableResources: true, } + mlog.DisableZap() + mainHelper = testlib.NewMainHelperWithOptions(&options) defer mainHelper.Close() diff --git a/cmd/mattermost/commands/main_test.go b/cmd/mattermost/commands/main_test.go index eb0bcfe772..3fae93f049 100644 --- a/cmd/mattermost/commands/main_test.go +++ b/cmd/mattermost/commands/main_test.go @@ -9,6 +9,7 @@ import ( "testing" "github.com/mattermost/mattermost-server/v5/api4" + "github.com/mattermost/mattermost-server/v5/mlog" "github.com/mattermost/mattermost-server/v5/testlib" ) @@ -27,6 +28,8 @@ func TestMain(m *testing.M) { EnableResources: true, } + mlog.DisableZap() + mainHelper = testlib.NewMainHelperWithOptions(&options) defer mainHelper.Close() api4.SetMainHelper(mainHelper) diff --git a/config/main_test.go b/config/main_test.go index a7cf254b8a..1b03e9c31e 100644 --- a/config/main_test.go +++ b/config/main_test.go @@ -9,6 +9,7 @@ import ( "github.com/go-sql-driver/mysql" "github.com/lib/pq" + "github.com/mattermost/mattermost-server/v5/mlog" "github.com/mattermost/mattermost-server/v5/model" "github.com/mattermost/mattermost-server/v5/testlib" "github.com/stretchr/testify/require" @@ -21,6 +22,8 @@ func TestMain(m *testing.M) { EnableStore: true, } + mlog.DisableZap() + mainHelper = testlib.NewMainHelperWithOptions(&options) defer mainHelper.Close() diff --git a/go.mod b/go.mod index 6e01325508..66c16af74b 100644 --- a/go.mod +++ b/go.mod @@ -60,7 +60,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.9 + github.com/mattermost/logr v1.0.13 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 f2686791ac..0e3f3edac4 100644 --- a/go.sum +++ b/go.sum @@ -484,6 +484,14 @@ github.com/mattermost/ldap v0.0.0-20191128190019-9f62ba4b8d4d h1:2DV7VIlEv6J5R5o github.com/mattermost/ldap v0.0.0-20191128190019-9f62ba4b8d4d/go.mod h1:HLbgMEI5K131jpxGazJ97AxfPDt31osq36YS1oxFQPQ= github.com/mattermost/logr v1.0.9 h1:jw6f6CjPC2YfPqzpGVM/vCGuqKLJdVS400ZRTFtEQwQ= github.com/mattermost/logr v1.0.9/go.mod h1:Mt4DPu1NXMe6JxPdwCC0XBoxXmN9eXOIRPoZarU2PXs= +github.com/mattermost/logr v1.0.10 h1:J/M6OFJhzQCUPGLyL9s8hiE+8nyL7Y0DybbOxYOisi0= +github.com/mattermost/logr v1.0.10/go.mod h1:Mt4DPu1NXMe6JxPdwCC0XBoxXmN9eXOIRPoZarU2PXs= +github.com/mattermost/logr v1.0.11 h1:XlNLB3x9OhvoNxEus46dW38zejINe5D2dZBlwmIlX5Q= +github.com/mattermost/logr v1.0.11/go.mod h1:Mt4DPu1NXMe6JxPdwCC0XBoxXmN9eXOIRPoZarU2PXs= +github.com/mattermost/logr v1.0.12 h1:1Tt2dJppjW6XlpJgMpeN+SNG1QgbTr4ITYnxG3NLPbM= +github.com/mattermost/logr v1.0.12/go.mod h1:Mt4DPu1NXMe6JxPdwCC0XBoxXmN9eXOIRPoZarU2PXs= +github.com/mattermost/logr v1.0.13 h1:6F/fM3csvH6Oy5sUpJuW7YyZSzZZAhJm5VcgKMxA2P8= +github.com/mattermost/logr v1.0.13/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/migrations/main_test.go b/migrations/main_test.go index 39c633b2a8..cffca24439 100644 --- a/migrations/main_test.go +++ b/migrations/main_test.go @@ -6,6 +6,7 @@ package migrations import ( "testing" + "github.com/mattermost/mattermost-server/v5/mlog" "github.com/mattermost/mattermost-server/v5/testlib" ) @@ -17,6 +18,8 @@ func TestMain(m *testing.M) { EnableResources: true, } + mlog.DisableZap() + mainHelper = testlib.NewMainHelperWithOptions(&options) defer mainHelper.Close() diff --git a/mlog/default.go b/mlog/default.go index 6322650d6b..1e409b192c 100644 --- a/mlog/default.go +++ b/mlog/default.go @@ -76,12 +76,18 @@ func defaultAdvancedShutdown(ctx context.Context) error { return nil } -func defaultAddTarget(target logr.Target) error { +func defaultAddTarget(targets ...logr.Target) error { // mlog.AddTarget should not be called until default // logger is replaced with mlog.Logger instance. return errors.New("cannot AddTarget on default logger") } +func defaultRemoveTargets(ctx context.Context, f func(TargetInfo) bool) error { + // mlog.RemoveTargets should not be called until default + // logger is replaced with mlog.Logger instance. + return errors.New("cannot RemoveTargets on default logger") +} + func defaultEnableMetrics(collector logr.MetricsCollector) error { // mlog.EnableMetrics should not be called until default // logger is replaced with mlog.Logger instance. diff --git a/mlog/global.go b/mlog/global.go index 39b4b30e2f..2986d92d29 100644 --- a/mlog/global.go +++ b/mlog/global.go @@ -5,6 +5,8 @@ package mlog import ( "context" + "log" + "sync/atomic" "github.com/mattermost/logr" "go.uber.org/zap" @@ -32,11 +34,29 @@ func InitGlobalLogger(logger *Logger) { ConfigAdvancedLogging = globalLogger.ConfigAdvancedLogging ShutdownAdvancedLogging = globalLogger.ShutdownAdvancedLogging AddTarget = globalLogger.AddTarget + RemoveTargets = globalLogger.RemoveTargets EnableMetrics = globalLogger.EnableMetrics } +// logWriterFunc provides access to mlog via io.Writer, so the standard logger +// can be redirected to use mlog and whatever targets are defined. +type logWriterFunc func([]byte) (int, error) + +func (lw logWriterFunc) Write(p []byte) (int, error) { + return lw(p) +} + func RedirectStdLog(logger *Logger) { - zap.RedirectStdLogAt(logger.zap.With(zap.String("source", "stdlog")).WithOptions(zap.AddCallerSkip(-2)), zapcore.ErrorLevel) + if atomic.LoadInt32(&disableZap) == 0 { + zap.RedirectStdLogAt(logger.zap.With(zap.String("source", "stdlog")).WithOptions(zap.AddCallerSkip(-2)), zapcore.ErrorLevel) + return + } + + writer := func(p []byte) (int, error) { + Log(LvlStdLog, string(p)) + return len(p), nil + } + log.SetOutput(logWriterFunc(writer)) } type LogFunc func(string, ...Field) @@ -45,7 +65,8 @@ type LogFuncCustomMulti func([]LogLevel, string, ...Field) 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 AddTargetFunc func(...logr.Target) error +type RemoveTargetsFunc func(context.Context, func(TargetInfo) bool) error type EnableMetricsFunc func(logr.MetricsCollector) error // DON'T USE THIS Modify the level on the app logger @@ -70,4 +91,5 @@ var Flush FlushFunc = defaultFlush var ConfigAdvancedLogging ConfigFunc = defaultAdvancedConfig var ShutdownAdvancedLogging ShutdownFunc = defaultAdvancedShutdown var AddTarget AddTargetFunc = defaultAddTarget +var RemoveTargets RemoveTargetsFunc = defaultRemoveTargets var EnableMetrics EnableMetricsFunc = defaultEnableMetrics diff --git a/mlog/levels.go b/mlog/levels.go index 898e82eb52..54bd25496e 100644 --- a/mlog/levels.go +++ b/mlog/levels.go @@ -5,13 +5,15 @@ package mlog // Standard levels var ( - LvlPanic = LogLevel{ID: 0, Name: "panic"} - LvlFatal = LogLevel{ID: 1, Name: "fatal"} + LvlPanic = LogLevel{ID: 0, Name: "panic", Stacktrace: true} + LvlFatal = LogLevel{ID: 1, Name: "fatal", Stacktrace: true} LvlError = LogLevel{ID: 2, Name: "error"} LvlWarn = LogLevel{ID: 3, Name: "warn"} LvlInfo = LogLevel{ID: 4, Name: "info"} LvlDebug = LogLevel{ID: 5, Name: "debug"} LvlTrace = LogLevel{ID: 6, Name: "trace"} + // used by redirected standard logger + LvlStdLog = LogLevel{ID: 10, Name: "stdlog"} // used only by the logger LvlLogError = LogLevel{ID: 11, Name: "logerror", Stacktrace: true} ) diff --git a/mlog/log.go b/mlog/log.go index c314a5732f..eaa8c10948 100644 --- a/mlog/log.go +++ b/mlog/log.go @@ -5,9 +5,12 @@ package mlog import ( "context" + "fmt" "io" "log" "os" + "sync/atomic" + "time" "github.com/mattermost/logr" "go.uber.org/zap" @@ -24,6 +27,19 @@ const ( LevelWarn = "warn" // Errors are messages about things we know are problems LevelError = "error" + + // DefaultFlushTimeout is the default amount of time mlog.Flush will wait + // before timing out. + DefaultFlushTimeout = time.Second * 5 +) + +var ( + // disableZap is set when Zap should be disabled and Logr used instead. + // This is needed for unit testing as Zap has no shutdown capabilities + // and holds file handles until process exit. Currently unit test create + // many server instances, and thus many Zap log files. + // This flag will be removed when Zap is permanently replaced. + disableZap int32 ) // Type and function aliases from zap to limit the libraries scope into MM code @@ -40,6 +56,8 @@ var NamedErr = zap.NamedError var Bool = zap.Bool var Duration = zap.Duration +type TargetInfo logr.TargetInfo + type LoggerConfiguration struct { EnableConsole bool ConsoleJson bool @@ -97,14 +115,33 @@ func NewLogger(config *LoggerConfiguration) *Logger { } if config.EnableFile { - writer := zapcore.AddSync(&lumberjack.Logger{ - Filename: config.FileLocation, - MaxSize: 100, - Compress: true, - }) + if atomic.LoadInt32(&disableZap) != 0 { + t := &LogTarget{ + Type: "file", + Format: "json", + Levels: mlogLevelToLogrLevels(config.FileLevel), + MaxQueueSize: DefaultMaxTargetQueue, + Options: []byte(fmt.Sprintf(`{"Filename":"%s", "MaxSizeMB":%d, "Compress":%t}`, + config.FileLocation, 100, true)), + } + if !config.FileJson { + t.Format = "plain" + } + if tgt, err := NewLogrTarget("mlogFile", t); err == nil { + logger.logrLogger.Logr().AddTarget(tgt) + } else { + Error("error creating mlogFile", Err(err)) + } + } else { + writer := zapcore.AddSync(&lumberjack.Logger{ + Filename: config.FileLocation, + MaxSize: 100, + Compress: true, + }) - core := zapcore.NewCore(makeEncoder(config.FileJson), writer, logger.fileLevel) - cores = append(cores, core) + core := zapcore.NewCore(makeEncoder(config.FileJson), writer, logger.fileLevel) + cores = append(cores, core) + } } combinedCore := zapcore.NewTee(cores...) @@ -221,14 +258,14 @@ func (l *Logger) LogM(levels []LogLevel, message string, fields ...Field) { } func (l *Logger) Flush(cxt context.Context) error { - return l.logrLogger.Logr().Flush() // TODO: use context when Logr lib supports it. + return l.logrLogger.Logr().FlushWithTimeout(cxt) } // 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 { - err := l.logrLogger.Logr().Shutdown() // TODO: use context when Logr lib supports it. + err := l.logrLogger.Logr().ShutdownWithTimeout(cxt) l.logrLogger = newLogr() return err } @@ -245,11 +282,22 @@ func (l *Logger) ConfigAdvancedLogging(targets LogTargetCfg) error { return err } -// AddTarget adds a logr.Target to the advanced logger. This is the preferred method +// AddTarget adds one or more logr.Target to the advanced logger. This is the preferred method // to add custom targets or provide configuration that cannot be expressed via a -//config source. -func (l *Logger) AddTarget(target logr.Target) error { - return l.logrLogger.Logr().AddTarget(target) +// config source. +func (l *Logger) AddTarget(targets ...logr.Target) error { + return l.logrLogger.Logr().AddTarget(targets...) +} + +// RemoveTargets selectively removes targets that were previously added to this logger instance +// using the passed in filter function. The filter function should return true to remove the target +// and false to keep it. +func (l *Logger) RemoveTargets(ctx context.Context, f func(ti TargetInfo) bool) error { + // Use locally defined TargetInfo type so we don't spread Logr dependencies. + fc := func(tic logr.TargetInfo) bool { + return f(TargetInfo(tic)) + } + return l.logrLogger.Logr().RemoveTargets(ctx, fc) } // EnableMetrics enables metrics collection by supplying a MetricsCollector. @@ -257,3 +305,21 @@ func (l *Logger) AddTarget(target logr.Target) error { func (l *Logger) EnableMetrics(collector logr.MetricsCollector) error { return l.logrLogger.Logr().SetMetricsCollector(collector) } + +// DisableZap is called to disable Zap, and Logr will be used instead. Any Logger +// instances created after this call will only use Logr. +// +// This is needed for unit testing as Zap has no shutdown capabilities +// and holds file handles until process exit. Currently unit tests create +// many server instances, and thus many Zap log file handles. +// +// This method will be removed when Zap is permanently replaced. +func DisableZap() { + atomic.StoreInt32(&disableZap, 1) +} + +// EnableZap re-enables Zap such that any Logger instances created after this +// call will allow Zap targets. +func EnableZap() { + atomic.StoreInt32(&disableZap, 0) +} diff --git a/mlog/logr.go b/mlog/logr.go index dba08e1006..01b39024aa 100644 --- a/mlog/logr.go +++ b/mlog/logr.go @@ -60,6 +60,7 @@ func logrAddTargets(logger *logr.Logger, targets LogTargetCfg) error { continue } if target != nil { + target.SetName(name) lgr.AddTarget(target) } } @@ -152,15 +153,18 @@ func newFileTarget(name string, t *LogTarget, filter logr.Filter, formatter logr if err := json.Unmarshal(t.Options, options); err != nil { return nil, err } + return newFileTargetWithOpts(name, t, target.FileOptions(*options), filter, formatter) +} - if options.Filename == "" { +func newFileTargetWithOpts(name string, t *LogTarget, opts target.FileOptions, filter logr.Filter, formatter logr.Formatter) (logr.Target, error) { + if opts.Filename == "" { return nil, fmt.Errorf("missing 'Filename' option for target %s", name) } - if err := checkFileWritable(options.Filename); err != nil { + if err := checkFileWritable(opts.Filename); err != nil { return nil, fmt.Errorf("error writing to 'Filename' for target %s: %w", name, err) } - newTarget := target.NewFileTarget(filter, formatter, target.FileOptions(*options), t.MaxQueueSize) + newTarget := target.NewFileTarget(filter, formatter, opts, t.MaxQueueSize) return newTarget, nil } @@ -222,3 +226,22 @@ func zapToLogr(zapFields []Field) logr.Fields { } return logr.Fields(encoder.Fields) } + +// mlogLevelToLogrLevel converts a mlog logger level to +// an array of discrete Logr levels. +func mlogLevelToLogrLevels(level string) []LogLevel { + levels := make([]LogLevel, 0) + levels = append(levels, LvlError, LvlPanic, LvlFatal, LvlStdLog) + + switch level { + case LevelDebug: + levels = append(levels, LvlDebug) + fallthrough + case LevelInfo: + levels = append(levels, LvlInfo) + fallthrough + case LevelWarn: + levels = append(levels, LvlWarn) + } + return levels +} diff --git a/services/slackimport/main_test.go b/services/slackimport/main_test.go index 28e21ea7ea..86bf3cba94 100644 --- a/services/slackimport/main_test.go +++ b/services/slackimport/main_test.go @@ -7,6 +7,8 @@ import ( "fmt" "os" "testing" + + "github.com/mattermost/mattermost-server/v5/mlog" ) func TestMain(m *testing.M) { @@ -20,6 +22,8 @@ func TestMain(m *testing.M) { panic(fmt.Sprintf("Failed to set current working directory to %s: %s", "../..", err.Error())) } + mlog.DisableZap() + defer func() { err := os.Chdir(prevDir) if err != nil { diff --git a/store/localcachelayer/main_test.go b/store/localcachelayer/main_test.go index fafd171529..bbbacba0f5 100644 --- a/store/localcachelayer/main_test.go +++ b/store/localcachelayer/main_test.go @@ -7,6 +7,7 @@ import ( "fmt" "testing" + "github.com/mattermost/mattermost-server/v5/mlog" "github.com/mattermost/mattermost-server/v5/services/cache" "github.com/mattermost/mattermost-server/v5/model" @@ -153,6 +154,7 @@ func getMockStore() *mocks.Store { } func TestMain(m *testing.M) { + mlog.DisableZap() mainHelper = testlib.NewMainHelperWithOptions(nil) defer mainHelper.Close() diff --git a/store/sqlstore/main_test.go b/store/sqlstore/main_test.go index 1c639241f7..c67cb02186 100644 --- a/store/sqlstore/main_test.go +++ b/store/sqlstore/main_test.go @@ -4,15 +4,18 @@ package sqlstore_test import ( - "github.com/mattermost/mattermost-server/v5/store/sqlstore" "testing" + "github.com/mattermost/mattermost-server/v5/mlog" + "github.com/mattermost/mattermost-server/v5/store/sqlstore" + "github.com/mattermost/mattermost-server/v5/testlib" ) var mainHelper *testlib.MainHelper func TestMain(m *testing.M) { + mlog.DisableZap() mainHelper = testlib.NewMainHelperWithOptions(nil) defer mainHelper.Close() diff --git a/vendor/github.com/mattermost/logr/logr.go b/vendor/github.com/mattermost/logr/logr.go index 9cda6d1cf0..631366a570 100644 --- a/vendor/github.com/mattermost/logr/logr.go +++ b/vendor/github.com/mattermost/logr/logr.go @@ -27,12 +27,13 @@ type Logr struct { shutdown bool lvlCache levelCache - metricsOnce sync.Once - metricsDone chan struct{} - metrics MetricsCollector - queueSizeGauge Gauge - loggedCounter Counter - errorCounter Counter + metricsInitOnce sync.Once + metricsCloseOnce sync.Once + metricsDone chan struct{} + metrics MetricsCollector + queueSizeGauge Gauge + loggedCounter Counter + errorCounter Counter bufferPool sync.Pool @@ -114,42 +115,33 @@ func (logr *Logr) Configure(config *cfg.Config) error { return fmt.Errorf("not implemented yet") } -// AddTarget adds a target to the logger which will receive -// log records for outputting. -func (logr *Logr) AddTarget(target Target) error { - logr.mux.Lock() - defer logr.mux.Unlock() - - if logr.shutdown { - return fmt.Errorf("logr shut down") - } - - logr.tmux.Lock() - 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) - } - } - +func (logr *Logr) ensureInit() { logr.once.Do(func() { + defer func() { + go logr.start() + }() + + logr.mux.Lock() + defer logr.mux.Unlock() + logr.maxQueueSizeActual = logr.MaxQueueSize if logr.maxQueueSizeActual == 0 { logr.maxQueueSizeActual = DefaultMaxQueueSize } + if logr.maxQueueSizeActual < 0 { logr.maxQueueSizeActual = 0 } + logr.in = make(chan *LogRec, logr.maxQueueSizeActual) logr.done = make(chan struct{}) + if logr.UseSyncMapLevelCache { logr.lvlCache = &syncMapLevelCache{} } else { logr.lvlCache = &arrayLevelCache{} } + if logr.MaxPooledBuffer == 0 { logr.MaxPooledBuffer = DefaultMaxPooledBuffer } @@ -158,11 +150,41 @@ func (logr *Logr) AddTarget(target Target) error { return new(bytes.Buffer) }, } + logr.lvlCache.setup() - go logr.start() }) - logr.resetLevelCache() - return err +} + +// AddTarget adds one or more targets to the logger which will receive +// log records for outputting. +func (logr *Logr) AddTarget(targets ...Target) error { + if logr.IsShutdown() { + return fmt.Errorf("AddTarget called after Logr shut down") + } + + logr.ensureInit() + metrics := logr.getMetricsCollector() + defer logr.ResetLevelCache() // call this after tmux is released + + logr.tmux.Lock() + defer logr.tmux.Unlock() + + errs := merror.New() + for _, t := range targets { + if t == nil { + continue + } + + logr.targets = append(logr.targets, t) + if metrics != nil { + if tm, ok := t.(TargetWithMetrics); ok { + if err := tm.EnableMetrics(metrics, logr.MetricsUpdateFreqMillis); err != nil { + errs.Append(err) + } + } + } + } + return errs.ErrorOrNil() } // NewLogger creates a Logger using defaults. A `Logger` is light-weight @@ -178,28 +200,13 @@ var levelStatusDisabled = LevelStatus{} // IsLevelEnabled returns true if at least one target has the specified // level enabled. The result is cached so that subsequent checks are fast. func (logr *Logr) IsLevelEnabled(lvl Level) LevelStatus { - // Check cache. lvlCache may still be nil if no targets added. - if logr.lvlCache == nil { - return levelStatusDisabled - } - status, ok := logr.lvlCache.get(lvl.ID) + status, ok := logr.isLevelEnabledFromCache(lvl) if ok { return status } - logr.mux.RLock() - defer logr.mux.RUnlock() - - // Don't accept new log records after shutdown. - if logr.shutdown { - return levelStatusDisabled - } - - status = LevelStatus{} - // Check each target. logr.tmux.RLock() - defer logr.tmux.RUnlock() for _, t := range logr.targets { e, s := t.IsLevelEnabled(lvl) if e { @@ -210,15 +217,45 @@ func (logr *Logr) IsLevelEnabled(lvl Level) LevelStatus { } } } + logr.tmux.RUnlock() // Cache and return the result. - if err := logr.lvlCache.put(lvl.ID, status); err != nil { + if err := logr.updateLevelCache(lvl.ID, status); err != nil { logr.ReportError(err) return LevelStatus{} } return status } +func (logr *Logr) isLevelEnabledFromCache(lvl Level) (LevelStatus, bool) { + logr.mux.RLock() + defer logr.mux.RUnlock() + + // Don't accept new log records after shutdown. + if logr.shutdown { + return levelStatusDisabled, true + } + + // Check cache. lvlCache may still be nil if no targets added. + if logr.lvlCache == nil { + return levelStatusDisabled, true + } + status, ok := logr.lvlCache.get(lvl.ID) + if ok { + return status, true + } + return LevelStatus{}, false +} + +func (logr *Logr) updateLevelCache(id LevelID, status LevelStatus) error { + logr.mux.RLock() + defer logr.mux.RUnlock() + if logr.lvlCache != nil { + return logr.lvlCache.put(id, status) + } + return nil +} + // HasTargets returns true only if at least one target exists within the Logr. func (logr *Logr) HasTargets() bool { logr.tmux.RLock() @@ -226,6 +263,71 @@ func (logr *Logr) HasTargets() bool { return len(logr.targets) > 0 } +// TargetInfo provides name and type for a Target. +type TargetInfo struct { + Name string + Type string +} + +// TargetInfos enumerates all the targets added to this Logr. +// The resulting slice represents a snapshot at time of calling. +func (logr *Logr) TargetInfos() []TargetInfo { + logr.tmux.RLock() + defer logr.tmux.RUnlock() + + infos := make([]TargetInfo, 0) + + for _, t := range logr.targets { + inf := TargetInfo{ + Name: fmt.Sprintf("%v", t), + Type: fmt.Sprintf("%T", t), + } + infos = append(infos, inf) + } + return infos +} + +// RemoveTargets safely removes one or more targets based on the filtering method. +// f should return true to delete the target, false to keep it. +// When removing a target, best effort is made to write any queued log records before +// closing, with cxt determining how much time can be spent in total. +// Note, keep the timeout short since this method blocks certain logging operations. +func (logr *Logr) RemoveTargets(cxt context.Context, f func(ti TargetInfo) bool) error { + var removed bool + defer func() { + if removed { + // call this after tmux is released since + // it will lock mux and we don't want to + // introduce possible deadlock. + logr.ResetLevelCache() + } + }() + + errs := merror.New() + + logr.tmux.Lock() + defer logr.tmux.Unlock() + + cp := make([]Target, 0) + + for _, t := range logr.targets { + inf := TargetInfo{ + Name: fmt.Sprintf("%v", t), + Type: fmt.Sprintf("%T", t), + } + if f(inf) { + if err := t.Shutdown(cxt); err != nil { + errs.Append(err) + } + removed = true + } else { + cp = append(cp, t) + } + } + logr.targets = cp + return errs.ErrorOrNil() +} + // 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() { @@ -304,15 +406,24 @@ 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 { + ctx, cancel := context.WithTimeout(context.Background(), logr.flushTimeout()) + defer cancel() + return logr.FlushWithTimeout(ctx) +} + +// Flush blocks while flushing the logr queue and all target queues, by +// writing existing log records to valid targets. +// Any attempts to add new log records will block until flush is complete. +// Use `IsTimeoutError` to determine if the returned error is +// due to a timeout. +func (logr *Logr) FlushWithTimeout(ctx context.Context) error { if !logr.HasTargets() { return nil } - logr.mux.Lock() - defer logr.mux.Unlock() - - ctx, cancel := context.WithTimeout(context.Background(), logr.flushTimeout()) - defer cancel() + if logr.IsShutdown() { + return errors.New("Flush called on shut down Logr") + } rec := newFlushLogRec(logr.NewLogger()) logr.enqueue(rec) @@ -325,6 +436,15 @@ func (logr *Logr) Flush() error { return nil } +// IsShutdown returns true if this Logr instance has been shut down. +// No further log records can be enqueued and no targets added after +// shutdown. +func (logr *Logr) IsShutdown() bool { + logr.mux.Lock() + defer logr.mux.Unlock() + return logr.shutdown +} + // Shutdown cleanly stops the logging engine after making best efforts // to flush all targets. Call this function right before application // exit - logr cannot be restarted once shut down. @@ -332,6 +452,17 @@ func (logr *Logr) Flush() error { // timing out. Use `IsTimeoutError` to determine if the returned error is // due to a timeout. func (logr *Logr) Shutdown() error { + ctx, cancel := context.WithTimeout(context.Background(), logr.shutdownTimeout()) + defer cancel() + return logr.ShutdownWithTimeout(ctx) +} + +// Shutdown cleanly stops the logging engine after making best efforts +// to flush all targets. Call this function right before application +// exit - logr cannot be restarted once shut down. +// Use `IsTimeoutError` to determine if the returned error is due to a +// timeout. +func (logr *Logr) ShutdownWithTimeout(ctx context.Context) error { logr.mux.Lock() if logr.shutdown { logr.mux.Unlock() @@ -339,16 +470,15 @@ 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() + logr.metricsCloseOnce.Do(func() { + if logr.metricsDone != nil { + close(logr.metricsDone) + } + }) - ctx, cancel := context.WithTimeout(context.Background(), logr.shutdownTimeout()) - defer cancel() + errs := merror.New() // close the incoming channel and wait for read loop to exit. if logr.in != nil { @@ -377,9 +507,8 @@ 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() - } + logr.incErrorCounter() + if logr.OnLoggerError == nil { fmt.Fprintln(os.Stderr, err) return @@ -455,7 +584,7 @@ func (logr *Logr) start() { // logr is closed. func (logr *Logr) startMetricsUpdater() { for { - updateFreq := logr.MetricsUpdateFreqMillis + updateFreq := logr.getMetricsUpdateFreqMillis() if updateFreq == 0 { updateFreq = DefMetricsUpdateFreqMillis } @@ -467,13 +596,17 @@ func (logr *Logr) startMetricsUpdater() { case <-logr.metricsDone: return case <-time.After(time.Duration(updateFreq) * time.Millisecond): - if logr.queueSizeGauge != nil { - logr.queueSizeGauge.Set(float64(len(logr.in))) - } + logr.setQueueSizeGauge(float64(len(logr.in))) } } } +func (logr *Logr) getMetricsUpdateFreqMillis() int64 { + logr.mux.RLock() + defer logr.mux.RUnlock() + return logr.MetricsUpdateFreqMillis +} + // fanout pushes a LogRec to all targets. func (logr *Logr) fanout(rec *LogRec) { var target Target @@ -484,6 +617,11 @@ func (logr *Logr) fanout(rec *LogRec) { }() var logged bool + defer func() { + if logged { + logr.incLoggedCounter() // call this after tmux is released + } + }() logr.tmux.RLock() defer logr.tmux.RUnlock() @@ -493,10 +631,6 @@ func (logr *Logr) fanout(rec *LogRec) { 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 index ac992b243b..24fe22b6e5 100644 --- a/vendor/github.com/mattermost/logr/metrics.go +++ b/vendor/github.com/mattermost/logr/metrics.go @@ -52,6 +52,12 @@ type TargetWithMetrics interface { EnableMetrics(collector MetricsCollector, updateFreqMillis int64) error } +func (logr *Logr) getMetricsCollector() MetricsCollector { + logr.mux.RLock() + defer logr.mux.RUnlock() + return logr.metrics +} + // 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 { @@ -59,12 +65,14 @@ func (logr *Logr) SetMetricsCollector(collector MetricsCollector) error { return errors.New("collector cannot be nil") } + logr.mux.Lock() logr.metrics = collector logr.queueSizeGauge, _ = collector.QueueSizeGauge("_logr") logr.loggedCounter, _ = collector.LoggedCounter("_logr") logr.errorCounter, _ = collector.ErrorCounter("_logr") + logr.mux.Unlock() - logr.metricsOnce.Do(func() { + logr.metricsInitOnce.Do(func() { logr.metricsDone = make(chan struct{}) go logr.startMetricsUpdater() }) @@ -75,7 +83,7 @@ func (logr *Logr) SetMetricsCollector(collector MetricsCollector) error { 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 { + if err := tm.EnableMetrics(collector, logr.MetricsUpdateFreqMillis); err != nil { merr.Append(err) } } @@ -83,3 +91,27 @@ func (logr *Logr) SetMetricsCollector(collector MetricsCollector) error { } return merr.ErrorOrNil() } + +func (logr *Logr) setQueueSizeGauge(val float64) { + logr.mux.RLock() + defer logr.mux.RUnlock() + if logr.queueSizeGauge != nil { + logr.queueSizeGauge.Set(val) + } +} + +func (logr *Logr) incLoggedCounter() { + logr.mux.RLock() + defer logr.mux.RUnlock() + if logr.loggedCounter != nil { + logr.loggedCounter.Inc() + } +} + +func (logr *Logr) incErrorCounter() { + logr.mux.RLock() + defer logr.mux.RUnlock() + if logr.errorCounter != nil { + logr.errorCounter.Inc() + } +} diff --git a/vendor/github.com/mattermost/logr/target.go b/vendor/github.com/mattermost/logr/target.go index 2ce76333cf..f8e7bf752c 100644 --- a/vendor/github.com/mattermost/logr/target.go +++ b/vendor/github.com/mattermost/logr/target.go @@ -4,13 +4,14 @@ import ( "context" "fmt" "os" + "sync" "time" ) // 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 provides an optional name for the target. SetName(name string) // IsLevelEnabled returns true if this target should emit @@ -39,7 +40,6 @@ type RecordWriter interface { // in your target type, implement `RecordWriter`, and call `(*Basic).Start`. type Basic struct { target Target - name string filter Filter formatter Formatter @@ -48,6 +48,10 @@ type Basic struct { done chan struct{} w RecordWriter + mux sync.RWMutex + name string + + metrics bool queueSizeGauge Gauge loggedCounter Counter errorCounter Counter @@ -74,12 +78,14 @@ func (b *Basic) Start(target Target, rw RecordWriter, filter Filter, formatter F b.w = rw go b.start() - if b.queueSizeGauge != nil { + if b.hasMetrics() { go b.startMetricsUpdater() } } func (b *Basic) SetName(name string) { + b.mux.Lock() + defer b.mux.Unlock() b.name = name } @@ -117,14 +123,10 @@ 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() - } + b.incDroppedCounter() return // drop the record } - if b.blockedCounter != nil { - b.blockedCounter.Inc() - } + b.incBlockedCounter() select { case <-time.After(lgr.enqueueTimeout()): @@ -136,9 +138,14 @@ func (b *Basic) Log(rec *LogRec) { // Metrics enables metrics collection using the provided MetricsCollector. func (b *Basic) EnableMetrics(collector MetricsCollector, updateFreqMillis int64) error { + name := fmt.Sprintf("%v", b) + + b.mux.Lock() + defer b.mux.Unlock() + + b.metrics = true b.metricsUpdateFreqMillis = updateFreqMillis - name := fmt.Sprintf("%v", b) var err error if b.queueSizeGauge, err = collector.QueueSizeGauge(name); err != nil { @@ -159,8 +166,57 @@ func (b *Basic) EnableMetrics(collector MetricsCollector, updateFreqMillis int64 return nil } +func (b *Basic) hasMetrics() bool { + b.mux.RLock() + defer b.mux.RUnlock() + return b.metrics +} + +func (b *Basic) setQueueSizeGauge(val float64) { + b.mux.RLock() + defer b.mux.RUnlock() + if b.queueSizeGauge != nil { + b.queueSizeGauge.Set(val) + } +} + +func (b *Basic) incLoggedCounter() { + b.mux.RLock() + defer b.mux.RUnlock() + if b.loggedCounter != nil { + b.loggedCounter.Inc() + } +} + +func (b *Basic) incErrorCounter() { + b.mux.RLock() + defer b.mux.RUnlock() + if b.errorCounter != nil { + b.errorCounter.Inc() + } +} + +func (b *Basic) incDroppedCounter() { + b.mux.RLock() + defer b.mux.RUnlock() + if b.droppedCounter != nil { + b.droppedCounter.Inc() + } +} + +func (b *Basic) incBlockedCounter() { + b.mux.RLock() + defer b.mux.RUnlock() + if b.blockedCounter != nil { + b.blockedCounter.Inc() + } +} + // String returns a name for this target. Use `SetName` to specify a name. func (b *Basic) String() string { + b.mux.RLock() + defer b.mux.RUnlock() + if b.name != "" { return b.name } @@ -183,12 +239,10 @@ func (b *Basic) start() { } else { err := b.w.Write(rec) if err != nil { - if b.errorCounter != nil { - b.errorCounter.Inc() - } + b.incErrorCounter() rec.Logger().Logr().ReportError(err) - } else if b.loggedCounter != nil { - b.loggedCounter.Inc() + } else { + b.incLoggedCounter() } } } @@ -199,7 +253,7 @@ func (b *Basic) start() { // target is closed. func (b *Basic) startMetricsUpdater() { for { - updateFreq := b.metricsUpdateFreqMillis + updateFreq := b.getMetricsUpdateFreqMillis() if updateFreq == 0 { updateFreq = DefMetricsUpdateFreqMillis } @@ -211,13 +265,17 @@ func (b *Basic) startMetricsUpdater() { case <-b.done: return case <-time.After(time.Duration(updateFreq) * time.Millisecond): - if b.queueSizeGauge != nil { - b.queueSizeGauge.Set(float64(len(b.in))) - } + b.setQueueSizeGauge(float64(len(b.in))) } } } +func (b *Basic) getMetricsUpdateFreqMillis() int64 { + b.mux.RLock() + defer b.mux.RUnlock() + return b.metricsUpdateFreqMillis +} + // flush drains the queue and notifies when done. func (b *Basic) flush(done chan<- struct{}) { for { @@ -229,9 +287,7 @@ 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() - } + b.incErrorCounter() rec.Logger().Logr().ReportError(err) } } diff --git a/vendor/modules.txt b/vendor/modules.txt index c1a8229bd3..29cc506951 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -281,7 +281,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.9 +# github.com/mattermost/logr v1.0.13 ## explicit github.com/mattermost/logr github.com/mattermost/logr/format diff --git a/web/main_test.go b/web/main_test.go index 7f47d6a1a0..fc08f7a5bd 100644 --- a/web/main_test.go +++ b/web/main_test.go @@ -6,6 +6,7 @@ package web import ( "testing" + "github.com/mattermost/mattermost-server/v5/mlog" "github.com/mattermost/mattermost-server/v5/testlib" ) @@ -17,6 +18,8 @@ func TestMain(m *testing.M) { EnableResources: true, } + mlog.DisableZap() + mainHelper = testlib.NewMainHelperWithOptions(&options) defer mainHelper.Close() diff --git a/web/web_test.go b/web/web_test.go index 5c888bb5c7..59a807c153 100644 --- a/web/web_test.go +++ b/web/web_test.go @@ -15,6 +15,7 @@ import ( "github.com/mattermost/mattermost-server/v5/app" "github.com/mattermost/mattermost-server/v5/config" + "github.com/mattermost/mattermost-server/v5/mlog" "github.com/mattermost/mattermost-server/v5/model" "github.com/mattermost/mattermost-server/v5/plugin" "github.com/mattermost/mattermost-server/v5/store" @@ -76,6 +77,8 @@ func setupTestHelper(t testing.TB, store store.Store, includeCacheLayer bool) *T options = append(options, app.ConfigStore(memoryStore)) options = append(options, app.StoreOverride(mainHelper.Store)) + mlog.DisableZap() + s, err := app.NewServer(options...) if err != nil { panic(err)