MM-36764 mlog refactor (#18118)
Refactor mlog - simplify mlog by removing redundant code - remove Zap dependency - update unit test helpers - update logging config - update auditing
Этот коммит содержится в:
@@ -4,26 +4,32 @@
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"context"
|
||||
"bytes"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"os"
|
||||
|
||||
"github.com/mattermost/logr"
|
||||
)
|
||||
|
||||
// defaultLog manually encodes the log to STDERR, providing a basic, default logging implementation
|
||||
// before mlog is fully configured.
|
||||
func defaultLog(level, msg string, fields ...Field) {
|
||||
func defaultLog(level Level, msg string, fields ...Field) {
|
||||
mFields := make(map[string]string)
|
||||
buf := &bytes.Buffer{}
|
||||
|
||||
for _, fld := range fields {
|
||||
buf.Reset()
|
||||
fld.ValueString(buf, shouldQuote)
|
||||
mFields[fld.Key] = buf.String()
|
||||
}
|
||||
|
||||
log := struct {
|
||||
Level string `json:"level"`
|
||||
Message string `json:"msg"`
|
||||
Fields []Field `json:"fields,omitempty"`
|
||||
Level string `json:"level"`
|
||||
Message string `json:"msg"`
|
||||
Fields map[string]string `json:"fields,omitempty"`
|
||||
}{
|
||||
level,
|
||||
level.Name,
|
||||
msg,
|
||||
fields,
|
||||
mFields,
|
||||
}
|
||||
|
||||
if b, err := json.Marshal(log); err != nil {
|
||||
@@ -33,67 +39,25 @@ func defaultLog(level, msg string, fields ...Field) {
|
||||
}
|
||||
}
|
||||
|
||||
func defaultIsLevelEnabled(level LogLevel) bool {
|
||||
func defaultIsLevelEnabled(level Level) bool {
|
||||
return true
|
||||
}
|
||||
|
||||
func defaultDebugLog(msg string, fields ...Field) {
|
||||
defaultLog("debug", msg, fields...)
|
||||
func defaultCustomMultiLog(lvl []Level, msg string, fields ...Field) {
|
||||
for _, level := range lvl {
|
||||
defaultLog(level, msg, fields...)
|
||||
}
|
||||
}
|
||||
|
||||
func defaultInfoLog(msg string, fields ...Field) {
|
||||
defaultLog("info", msg, fields...)
|
||||
}
|
||||
|
||||
func defaultWarnLog(msg string, fields ...Field) {
|
||||
defaultLog("warn", msg, fields...)
|
||||
}
|
||||
|
||||
func defaultErrorLog(msg string, fields ...Field) {
|
||||
defaultLog("error", msg, fields...)
|
||||
}
|
||||
|
||||
func defaultCriticalLog(msg string, fields ...Field) {
|
||||
// We map critical to error in zap, so be consistent.
|
||||
defaultLog("error", msg, fields...)
|
||||
}
|
||||
|
||||
func defaultCustomLog(lvl LogLevel, msg string, fields ...Field) {
|
||||
// custom log levels are only output once log targets are configured.
|
||||
}
|
||||
|
||||
func defaultCustomMultiLog(lvl []LogLevel, msg string, fields ...Field) {
|
||||
// custom log levels are only output once log targets are configured.
|
||||
}
|
||||
|
||||
func defaultFlush(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func defaultAdvancedConfig(cfg LogTargetCfg) error {
|
||||
// mlog.ConfigAdvancedConfig should not be called until default
|
||||
// logger is replaced with mlog.Logger instance.
|
||||
return errors.New("cannot config advanced logging on default logger")
|
||||
}
|
||||
|
||||
func defaultAdvancedShutdown(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
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.
|
||||
return errors.New("cannot EnableMetrics on default logger")
|
||||
// shouldQuote returns true if val contains any characters that require quotations.
|
||||
func shouldQuote(val string) bool {
|
||||
for _, c := range val {
|
||||
if !((c >= '0' && c <= '9') ||
|
||||
(c >= 'a' && c <= 'z') ||
|
||||
(c >= 'A' && c <= 'Z') ||
|
||||
c == '-' || c == '.' || c == '_' || c == '/' || c == '@' || c == '^' || c == '+') {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
@@ -1,32 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"github.com/mattermost/logr"
|
||||
)
|
||||
|
||||
// onLoggerError is called when the logging system encounters an error,
|
||||
// such as a target not able to write records. The targets will keep trying
|
||||
// however the error will be logged with a dedicated level that can be output
|
||||
// to a safe/always available target for monitoring or alerting.
|
||||
func onLoggerError(err error) {
|
||||
Log(LvlLogError, "advanced logging error", Err(err))
|
||||
}
|
||||
|
||||
// onQueueFull is called when the main logger queue is full, indicating the
|
||||
// volume and frequency of log record creation is too high for the queue size
|
||||
// and/or the target latencies.
|
||||
func onQueueFull(rec *logr.LogRec, maxQueueSize int) bool {
|
||||
Log(LvlLogError, "main queue full, dropping record", Any("rec", rec))
|
||||
return true // drop record
|
||||
}
|
||||
|
||||
// onTargetQueueFull is called when the main logger queue is full, indicating the
|
||||
// volume and frequency of log record creation is too high for the target's queue size
|
||||
// and/or the target latency.
|
||||
func onTargetQueueFull(target logr.Target, rec *logr.LogRec, maxQueueSize int) bool {
|
||||
Log(LvlLogError, "target queue full, dropping record", String("target", ""), Any("rec", rec))
|
||||
return true // drop record
|
||||
}
|
||||
@@ -4,95 +4,119 @@
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"sync/atomic"
|
||||
|
||||
"github.com/mattermost/logr"
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
"sync"
|
||||
)
|
||||
|
||||
var globalLogger *Logger
|
||||
var (
|
||||
globalLogger *Logger
|
||||
muxGlobalLogger sync.RWMutex
|
||||
)
|
||||
|
||||
func InitGlobalLogger(logger *Logger) {
|
||||
// Clean up previous instance.
|
||||
if globalLogger != nil && globalLogger.logrLogger != nil {
|
||||
globalLogger.logrLogger.Logr().Shutdown()
|
||||
muxGlobalLogger.Lock()
|
||||
defer muxGlobalLogger.Unlock()
|
||||
|
||||
globalLogger = logger
|
||||
}
|
||||
|
||||
func getGlobalLogger() *Logger {
|
||||
muxGlobalLogger.RLock()
|
||||
defer muxGlobalLogger.RUnlock()
|
||||
|
||||
return globalLogger
|
||||
}
|
||||
|
||||
// IsLevelEnabled returns true only if at least one log target is
|
||||
// configured to emit the specified log level. Use this check when
|
||||
// gathering the log info may be expensive.
|
||||
//
|
||||
// Note, transformations and serializations done via fields are already
|
||||
// lazily evaluated and don't require this check beforehand.
|
||||
func IsLevelEnabled(level Level) bool {
|
||||
logger := getGlobalLogger()
|
||||
if logger == nil {
|
||||
return defaultIsLevelEnabled(level)
|
||||
}
|
||||
glob := *logger
|
||||
glob.zap = glob.zap.WithOptions(zap.AddCallerSkip(1))
|
||||
globalLogger = &glob
|
||||
IsLevelEnabled = globalLogger.IsLevelEnabled
|
||||
Debug = globalLogger.Debug
|
||||
Info = globalLogger.Info
|
||||
Warn = globalLogger.Warn
|
||||
Error = globalLogger.Error
|
||||
Critical = globalLogger.Critical
|
||||
Log = globalLogger.Log
|
||||
LogM = globalLogger.LogM
|
||||
Flush = globalLogger.Flush
|
||||
ConfigAdvancedLogging = globalLogger.ConfigAdvancedLogging
|
||||
ShutdownAdvancedLogging = globalLogger.ShutdownAdvancedLogging
|
||||
AddTarget = globalLogger.AddTarget
|
||||
RemoveTargets = globalLogger.RemoveTargets
|
||||
EnableMetrics = globalLogger.EnableMetrics
|
||||
return logger.IsLevelEnabled(level)
|
||||
}
|
||||
|
||||
// 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) {
|
||||
if atomic.LoadInt32(&disableZap) == 0 {
|
||||
zap.RedirectStdLogAt(logger.zap.With(zap.String("source", "stdlog")).WithOptions(zap.AddCallerSkip(-2)), zapcore.ErrorLevel)
|
||||
// Log emits the log record for any targets configured for the specified level.
|
||||
func Log(level Level, msg string, fields ...Field) {
|
||||
logger := getGlobalLogger()
|
||||
if logger == nil {
|
||||
defaultLog(level, msg, fields...)
|
||||
return
|
||||
}
|
||||
logger.Log(level, msg, fields...)
|
||||
}
|
||||
|
||||
writer := func(p []byte) (int, error) {
|
||||
Log(LvlStdLog, string(p))
|
||||
return len(p), nil
|
||||
// LogM emits the log record for any targets configured for the specified levels.
|
||||
// Equivalent to calling `Log` once for each level.
|
||||
func LogM(levels []Level, msg string, fields ...Field) {
|
||||
logger := getGlobalLogger()
|
||||
if logger == nil {
|
||||
defaultCustomMultiLog(levels, msg, fields...)
|
||||
return
|
||||
}
|
||||
log.SetOutput(logWriterFunc(writer))
|
||||
logger.LogM(levels, msg, fields...)
|
||||
}
|
||||
|
||||
type IsLevelEnabledFunc func(LogLevel) bool
|
||||
type LogFunc func(string, ...Field)
|
||||
type LogFuncCustom func(LogLevel, string, ...Field)
|
||||
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 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
|
||||
func GloballyDisableDebugLogForTest() {
|
||||
globalLogger.consoleLevel.SetLevel(zapcore.ErrorLevel)
|
||||
// Convenience method equivalent to calling `Log` with the `Trace` level.
|
||||
func Trace(msg string, fields ...Field) {
|
||||
logger := getGlobalLogger()
|
||||
if logger == nil {
|
||||
defaultLog(LvlTrace, msg, fields...)
|
||||
return
|
||||
}
|
||||
logger.Trace(msg, fields...)
|
||||
}
|
||||
|
||||
// DON'T USE THIS Modify the level on the app logger
|
||||
func GloballyEnableDebugLogForTest() {
|
||||
globalLogger.consoleLevel.SetLevel(zapcore.DebugLevel)
|
||||
// Convenience method equivalent to calling `Log` with the `Debug` level.
|
||||
func Debug(msg string, fields ...Field) {
|
||||
logger := getGlobalLogger()
|
||||
if logger == nil {
|
||||
defaultLog(LvlDebug, msg, fields...)
|
||||
return
|
||||
}
|
||||
logger.Debug(msg, fields...)
|
||||
}
|
||||
|
||||
var IsLevelEnabled IsLevelEnabledFunc = defaultIsLevelEnabled
|
||||
var Debug LogFunc = defaultDebugLog
|
||||
var Info LogFunc = defaultInfoLog
|
||||
var Warn LogFunc = defaultWarnLog
|
||||
var Error LogFunc = defaultErrorLog
|
||||
var Critical LogFunc = defaultCriticalLog
|
||||
var Log LogFuncCustom = defaultCustomLog
|
||||
var LogM LogFuncCustomMulti = defaultCustomMultiLog
|
||||
var Flush FlushFunc = defaultFlush
|
||||
// Convenience method equivalent to calling `Log` with the `Info` level.
|
||||
func Info(msg string, fields ...Field) {
|
||||
logger := getGlobalLogger()
|
||||
if logger == nil {
|
||||
defaultLog(LvlInfo, msg, fields...)
|
||||
return
|
||||
}
|
||||
logger.Info(msg, fields...)
|
||||
}
|
||||
|
||||
var ConfigAdvancedLogging ConfigFunc = defaultAdvancedConfig
|
||||
var ShutdownAdvancedLogging ShutdownFunc = defaultAdvancedShutdown
|
||||
var AddTarget AddTargetFunc = defaultAddTarget
|
||||
var RemoveTargets RemoveTargetsFunc = defaultRemoveTargets
|
||||
var EnableMetrics EnableMetricsFunc = defaultEnableMetrics
|
||||
// Convenience method equivalent to calling `Log` with the `Warn` level.
|
||||
func Warn(msg string, fields ...Field) {
|
||||
logger := getGlobalLogger()
|
||||
if logger == nil {
|
||||
defaultLog(LvlWarn, msg, fields...)
|
||||
return
|
||||
}
|
||||
logger.Warn(msg, fields...)
|
||||
}
|
||||
|
||||
// Convenience method equivalent to calling `Log` with the `Error` level.
|
||||
func Error(msg string, fields ...Field) {
|
||||
logger := getGlobalLogger()
|
||||
if logger == nil {
|
||||
defaultLog(LvlError, msg, fields...)
|
||||
return
|
||||
}
|
||||
logger.Error(msg, fields...)
|
||||
}
|
||||
|
||||
// Convenience method equivalent to calling `Log` with the `Critical` level.
|
||||
func Critical(msg string, fields ...Field) {
|
||||
logger := getGlobalLogger()
|
||||
if logger == nil {
|
||||
defaultLog(LvlCritical, msg, fields...)
|
||||
return
|
||||
}
|
||||
logger.Critical(msg, fields...)
|
||||
}
|
||||
|
||||
@@ -4,6 +4,8 @@
|
||||
package mlog_test
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"path/filepath"
|
||||
@@ -29,83 +31,86 @@ func TestLoggingBeforeInitialized(t *testing.T) {
|
||||
|
||||
func TestLoggingAfterInitialized(t *testing.T) {
|
||||
testCases := []struct {
|
||||
Description string
|
||||
LoggerConfiguration *mlog.LoggerConfiguration
|
||||
ExpectedLogs []string
|
||||
description string
|
||||
cfg mlog.TargetCfg
|
||||
expectedLogs []string
|
||||
}{
|
||||
{
|
||||
"file logging, json, debug",
|
||||
&mlog.LoggerConfiguration{
|
||||
EnableConsole: false,
|
||||
EnableFile: true,
|
||||
FileJson: true,
|
||||
FileLevel: mlog.LevelDebug,
|
||||
mlog.TargetCfg{
|
||||
Type: "file",
|
||||
Format: "json",
|
||||
FormatOptions: json.RawMessage(`{"enable_caller":true}`),
|
||||
Levels: []mlog.Level{mlog.LvlCritical, mlog.LvlError, mlog.LvlWarn, mlog.LvlInfo, mlog.LvlDebug},
|
||||
},
|
||||
[]string{
|
||||
`{"level":"debug","ts":0,"caller":"mlog/global_test.go:0","msg":"real debug log"}`,
|
||||
`{"level":"info","ts":0,"caller":"mlog/global_test.go:0","msg":"real info log"}`,
|
||||
`{"level":"warn","ts":0,"caller":"mlog/global_test.go:0","msg":"real warning log"}`,
|
||||
`{"level":"error","ts":0,"caller":"mlog/global_test.go:0","msg":"real error log"}`,
|
||||
`{"level":"error","ts":0,"caller":"mlog/global_test.go:0","msg":"real critical log"}`,
|
||||
`{"timestamp":0,"level":"debug","msg":"real debug log","caller":"mlog/global_test.go:0"}`,
|
||||
`{"timestamp":0,"level":"info","msg":"real info log","caller":"mlog/global_test.go:0"}`,
|
||||
`{"timestamp":0,"level":"warn","msg":"real warning log","caller":"mlog/global_test.go:0"}`,
|
||||
`{"timestamp":0,"level":"error","msg":"real error log","caller":"mlog/global_test.go:0"}`,
|
||||
`{"timestamp":0,"level":"critical","msg":"real critical log","caller":"mlog/global_test.go:0"}`,
|
||||
},
|
||||
},
|
||||
{
|
||||
"file logging, json, error",
|
||||
&mlog.LoggerConfiguration{
|
||||
EnableConsole: false,
|
||||
EnableFile: true,
|
||||
FileJson: true,
|
||||
FileLevel: mlog.LevelError,
|
||||
mlog.TargetCfg{
|
||||
Type: "file",
|
||||
Format: "json",
|
||||
FormatOptions: json.RawMessage(`{"enable_caller":true}`),
|
||||
Levels: []mlog.Level{mlog.LvlCritical, mlog.LvlError},
|
||||
},
|
||||
[]string{
|
||||
`{"level":"error","ts":0,"caller":"mlog/global_test.go:0","msg":"real error log"}`,
|
||||
`{"level":"error","ts":0,"caller":"mlog/global_test.go:0","msg":"real critical log"}`,
|
||||
`{"timestamp":0,"level":"error","msg":"real error log","caller":"mlog/global_test.go:0"}`,
|
||||
`{"timestamp":0,"level":"critical","msg":"real critical log","caller":"mlog/global_test.go:0"}`,
|
||||
},
|
||||
},
|
||||
{
|
||||
"file logging, non-json, debug",
|
||||
&mlog.LoggerConfiguration{
|
||||
EnableConsole: false,
|
||||
EnableFile: true,
|
||||
FileJson: false,
|
||||
FileLevel: mlog.LevelDebug,
|
||||
mlog.TargetCfg{
|
||||
Type: "file",
|
||||
Format: "plain",
|
||||
FormatOptions: json.RawMessage(`{"delim":" | ", "enable_caller":true}`),
|
||||
Levels: []mlog.Level{mlog.LvlCritical, mlog.LvlError, mlog.LvlWarn, mlog.LvlInfo, mlog.LvlDebug},
|
||||
},
|
||||
[]string{
|
||||
`TIME debug mlog/global_test.go:0 real debug log`,
|
||||
`TIME info mlog/global_test.go:0 real info log`,
|
||||
`TIME warn mlog/global_test.go:0 real warning log`,
|
||||
`TIME error mlog/global_test.go:0 real error log`,
|
||||
`TIME error mlog/global_test.go:0 real critical log`,
|
||||
`debug | TIME | real debug log | caller="mlog/global_test.go:0"`,
|
||||
`info | TIME | real info log | caller="mlog/global_test.go:0"`,
|
||||
`warn | TIME | real warning log | caller="mlog/global_test.go:0"`,
|
||||
`error | TIME | real error log | caller="mlog/global_test.go:0"`,
|
||||
`critical | TIME | real critical log | caller="mlog/global_test.go:0"`,
|
||||
},
|
||||
},
|
||||
{
|
||||
"file logging, non-json, error",
|
||||
&mlog.LoggerConfiguration{
|
||||
EnableConsole: false,
|
||||
EnableFile: true,
|
||||
FileJson: false,
|
||||
FileLevel: mlog.LevelError,
|
||||
mlog.TargetCfg{
|
||||
Type: "file",
|
||||
Format: "plain",
|
||||
FormatOptions: json.RawMessage(`{"delim":" | ", "enable_caller":true}`),
|
||||
Levels: []mlog.Level{mlog.LvlCritical, mlog.LvlError},
|
||||
},
|
||||
[]string{
|
||||
`TIME error mlog/global_test.go:0 real error log`,
|
||||
`TIME error mlog/global_test.go:0 real critical log`,
|
||||
`error | TIME | real error log | caller="mlog/global_test.go:0"`,
|
||||
`critical | TIME | real critical log | caller="mlog/global_test.go:0"`,
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
for _, testCase := range testCases {
|
||||
t.Run(testCase.Description, func(t *testing.T) {
|
||||
t.Run(testCase.description, func(t *testing.T) {
|
||||
var filePath string
|
||||
if testCase.LoggerConfiguration.EnableFile {
|
||||
if testCase.cfg.Type == "file" {
|
||||
tempDir, err := ioutil.TempDir(os.TempDir(), "TestLoggingAfterInitialized")
|
||||
require.NoError(t, err)
|
||||
defer os.Remove(tempDir)
|
||||
|
||||
filePath = filepath.Join(tempDir, "file.log")
|
||||
testCase.LoggerConfiguration.FileLocation = filePath
|
||||
testCase.cfg.Options = json.RawMessage(fmt.Sprintf(`{"filename": "%s"}`, filePath))
|
||||
}
|
||||
|
||||
logger := mlog.NewLogger(testCase.LoggerConfiguration)
|
||||
logger, _ := mlog.NewLogger()
|
||||
err := logger.ConfigureTargets(map[string]mlog.TargetCfg{testCase.description: testCase.cfg})
|
||||
require.NoError(t, err)
|
||||
|
||||
mlog.InitGlobalLogger(logger)
|
||||
|
||||
mlog.Debug("real debug log")
|
||||
@@ -114,32 +119,26 @@ func TestLoggingAfterInitialized(t *testing.T) {
|
||||
mlog.Error("real error log")
|
||||
mlog.Critical("real critical log")
|
||||
|
||||
if testCase.LoggerConfiguration.EnableFile {
|
||||
logger.Shutdown()
|
||||
|
||||
if testCase.cfg.Type == "file" {
|
||||
logs, err := ioutil.ReadFile(filePath)
|
||||
require.NoError(t, err)
|
||||
|
||||
actual := strings.TrimSpace(string(logs))
|
||||
|
||||
if testCase.LoggerConfiguration.FileJson {
|
||||
reTs := regexp.MustCompile(`"ts":[0-9\.]+`)
|
||||
if testCase.cfg.Format == "json" {
|
||||
reTs := regexp.MustCompile(`"timestamp":"[0-9\.\-\:\sZ]+"`)
|
||||
reCaller := regexp.MustCompile(`"caller":"([^"]+):[0-9\.]+"`)
|
||||
actual = reTs.ReplaceAllString(actual, `"ts":0`)
|
||||
actual = reTs.ReplaceAllString(actual, `"timestamp":0`)
|
||||
actual = reCaller.ReplaceAllString(actual, `"caller":"$1:0"`)
|
||||
} else {
|
||||
actualRows := strings.Split(actual, "\n")
|
||||
for i, actualRow := range actualRows {
|
||||
actualFields := strings.Split(actualRow, "\t")
|
||||
if len(actualFields) > 3 {
|
||||
actualFields[0] = "TIME"
|
||||
reCaller := regexp.MustCompile(`([^"]+):[0-9\.]+`)
|
||||
actualFields[2] = reCaller.ReplaceAllString(actualFields[2], "$1:0")
|
||||
actualRows[i] = strings.Join(actualFields, "\t")
|
||||
}
|
||||
}
|
||||
|
||||
actual = strings.Join(actualRows, "\n")
|
||||
reTs := regexp.MustCompile(`\[\d\d\d\d-\d\d-\d\d\s[0-9\:\.\s\-Z]+\]`)
|
||||
reCaller := regexp.MustCompile(`caller="([^"]+):[0-9\.]+"`)
|
||||
actual = reTs.ReplaceAllString(actual, "TIME")
|
||||
actual = reCaller.ReplaceAllString(actual, `caller="$1:0"`)
|
||||
}
|
||||
require.ElementsMatch(t, testCase.ExpectedLogs, strings.Split(actual, "\n"))
|
||||
require.ElementsMatch(t, testCase.expectedLogs, strings.Split(actual, "\n"))
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
@@ -1,52 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package human
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/mattermost/mattermost-server/v6/shared/mlog"
|
||||
)
|
||||
|
||||
type LogEntry struct {
|
||||
Time time.Time
|
||||
Level string
|
||||
Message string
|
||||
Caller string
|
||||
Fields []mlog.Field
|
||||
}
|
||||
|
||||
// Provide default string representation. Used by SimpleWriter
|
||||
func (f LogEntry) String() string {
|
||||
var sb strings.Builder
|
||||
if !f.Time.IsZero() {
|
||||
sb.WriteString(f.Time.Format(time.RFC3339Nano))
|
||||
sb.WriteRune(' ')
|
||||
}
|
||||
if f.Level != "" {
|
||||
sb.WriteString(f.Level)
|
||||
sb.WriteRune(' ')
|
||||
}
|
||||
if f.Caller != "" {
|
||||
sb.WriteString(f.Caller)
|
||||
sb.WriteRune(' ')
|
||||
}
|
||||
for _, field := range f.Fields {
|
||||
sb.WriteString(field.Key)
|
||||
sb.WriteRune('=')
|
||||
sb.WriteString(fmt.Sprint(field.Interface))
|
||||
sb.WriteRune(' ')
|
||||
}
|
||||
if f.Message != "" {
|
||||
// If the message is multiple lines, start the whole message on a new line
|
||||
if strings.ContainsRune(f.Message, '\n') {
|
||||
sb.WriteRune('\n')
|
||||
}
|
||||
sb.WriteString(f.Message)
|
||||
}
|
||||
|
||||
return sb.String()
|
||||
}
|
||||
@@ -1,77 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package human
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
"time"
|
||||
|
||||
"github.com/sirupsen/logrus"
|
||||
)
|
||||
|
||||
type LogrusWriter struct {
|
||||
logger *logrus.Logger
|
||||
}
|
||||
|
||||
func (w *LogrusWriter) Write(e LogEntry) {
|
||||
if e.Level == "" {
|
||||
fmt.Fprintln(w.logger.Out, e.Message)
|
||||
return
|
||||
}
|
||||
|
||||
lvl, err := logrus.ParseLevel(e.Level)
|
||||
if err != nil {
|
||||
fmt.Fprintln(w.logger.Out, err)
|
||||
lvl = logrus.TraceLevel + 1 // will invoke Println
|
||||
}
|
||||
|
||||
logger := w.logger.WithTime(e.Time)
|
||||
|
||||
if e.Caller != "" {
|
||||
// logrus has a system of reporting the caller, but there's no easy way to override it
|
||||
logger = logger.WithField("caller", e.Caller)
|
||||
}
|
||||
|
||||
for _, field := range e.Fields {
|
||||
logger = logger.WithField(field.Key, field.Interface)
|
||||
}
|
||||
|
||||
switch lvl {
|
||||
case logrus.PanicLevel:
|
||||
// Prevent panic from causing us to exit
|
||||
defer func() {
|
||||
recover()
|
||||
}()
|
||||
logger.Panic(e.Message)
|
||||
case logrus.FatalLevel:
|
||||
logger.Fatal(e.Message)
|
||||
case logrus.ErrorLevel:
|
||||
logger.Error(e.Message)
|
||||
case logrus.WarnLevel:
|
||||
logger.Warn(e.Message)
|
||||
case logrus.InfoLevel:
|
||||
logger.Info(e.Message)
|
||||
case logrus.DebugLevel:
|
||||
logger.Debug(e.Message)
|
||||
case logrus.TraceLevel:
|
||||
logger.Trace(e.Message)
|
||||
default:
|
||||
logger.Println(e.Message)
|
||||
}
|
||||
}
|
||||
|
||||
func NewLogrusWriter(output io.Writer) *LogrusWriter {
|
||||
w := new(LogrusWriter)
|
||||
w.logger = logrus.New()
|
||||
w.logger.SetLevel(logrus.TraceLevel) // don't filter any logs
|
||||
w.logger.ExitFunc = func(int) {} // prevent Fatal from causing us to exit
|
||||
w.logger.SetReportCaller(false)
|
||||
w.logger.SetOutput(output)
|
||||
var tf logrus.TextFormatter
|
||||
tf.FullTimestamp = true
|
||||
tf.TimestampFormat = time.RFC3339Nano
|
||||
w.logger.SetFormatter(&tf)
|
||||
return w
|
||||
}
|
||||
@@ -1,181 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package human
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/mattermost/mattermost-server/v6/shared/mlog"
|
||||
)
|
||||
|
||||
func ParseLogMessage(msg string) LogEntry {
|
||||
result, err := parseLogMessage(msg)
|
||||
if err != nil {
|
||||
// If failed to parse, just output a LogEntry where all fields are blank, but Message is the original string
|
||||
var result2 LogEntry
|
||||
result2.Message = msg
|
||||
return result2
|
||||
}
|
||||
return result
|
||||
}
|
||||
|
||||
func parseLogMessage(msg string) (result LogEntry, err error) {
|
||||
|
||||
// Note: This implementation uses a custom json decoding loop.
|
||||
// The primary advantage of this versus decoding directly into a map is to
|
||||
// preserve the order of the fields. This can be simplified if we end up
|
||||
// having the formatter sort fields alphabetically (logrus does by default)
|
||||
|
||||
dec := json.NewDecoder(strings.NewReader(msg))
|
||||
|
||||
// look for an initial "{"
|
||||
token, err := dec.Token()
|
||||
if err != nil {
|
||||
return result, err
|
||||
}
|
||||
d, ok := token.(json.Delim)
|
||||
if !ok || d != '{' {
|
||||
return result, fmt.Errorf("input is not a JSON object, found: %v", token)
|
||||
}
|
||||
|
||||
// read all key-value pairs
|
||||
for dec.More() {
|
||||
key, err2 := dec.Token()
|
||||
if err2 != nil {
|
||||
return result, err2
|
||||
}
|
||||
skey, ok2 := key.(string)
|
||||
if !ok2 {
|
||||
return result, errors.New("key is not a value string")
|
||||
}
|
||||
if !dec.More() {
|
||||
return result, errors.New("missing value pair")
|
||||
}
|
||||
|
||||
switch skey {
|
||||
case "ts":
|
||||
var ts json.Number
|
||||
if err2 := dec.Decode(&ts); err2 != nil {
|
||||
return result, err2
|
||||
}
|
||||
timeVal, err2 := numberToTime(ts)
|
||||
if err2 != nil {
|
||||
return result, err2
|
||||
}
|
||||
result.Time = timeVal
|
||||
|
||||
case "level":
|
||||
s, err2 := decodeAsString(dec)
|
||||
if err2 != nil {
|
||||
return result, err2
|
||||
}
|
||||
result.Level = s
|
||||
|
||||
case "msg":
|
||||
s, err2 := decodeAsString(dec)
|
||||
if err2 != nil {
|
||||
return result, err2
|
||||
}
|
||||
result.Message = s
|
||||
|
||||
case "caller":
|
||||
s, err2 := decodeAsString(dec)
|
||||
if err2 != nil {
|
||||
return result, err2
|
||||
}
|
||||
result.Caller = s
|
||||
|
||||
default:
|
||||
var p interface{}
|
||||
if err2 := dec.Decode(&p); err2 != nil {
|
||||
return result, err2
|
||||
}
|
||||
var f mlog.Field
|
||||
f.Key = skey
|
||||
f.Interface = p
|
||||
result.Fields = append(result.Fields, f)
|
||||
}
|
||||
}
|
||||
|
||||
// read the "}"
|
||||
token, err = dec.Token()
|
||||
if err != nil {
|
||||
return result, err
|
||||
}
|
||||
d, ok = token.(json.Delim)
|
||||
if !ok || d != '}' {
|
||||
return result, fmt.Errorf("failed to read '}', read: %v", token)
|
||||
}
|
||||
|
||||
// make sure nothing else trailing
|
||||
if token, err := dec.Token(); err != io.EOF {
|
||||
return result, err
|
||||
} else if token != nil {
|
||||
return result, errors.New("found trailing data")
|
||||
}
|
||||
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// Translate a number into a time
|
||||
func numberToTime(v json.Number) (time.Time, error) {
|
||||
// Using floating point math to extract the nanoseconds leads to a time that doesn't exactly match the input
|
||||
// Instead, parse out the components from the string representation
|
||||
|
||||
var t time.Time
|
||||
|
||||
// First make sure it is a number...
|
||||
flt, err := v.Float64()
|
||||
if err != nil {
|
||||
return t, err
|
||||
}
|
||||
|
||||
s := v.String()
|
||||
|
||||
if strings.ContainsAny(s, "eE") {
|
||||
// input is in scientific notation. Convert to standard decimal notation
|
||||
s = strconv.FormatFloat(flt, 'f', -1, 64)
|
||||
}
|
||||
|
||||
// extract the seconds and nanoseconds separately
|
||||
var nanos, sec int64
|
||||
|
||||
parts := strings.SplitN(s, ".", 2)
|
||||
sec, err = strconv.ParseInt(parts[0], 10, 64)
|
||||
if err != nil {
|
||||
return t, err
|
||||
}
|
||||
|
||||
if len(parts) == 2 {
|
||||
nanosText := parts[1] + "000000000"
|
||||
nanosText = nanosText[:9]
|
||||
nanos, err = strconv.ParseInt(nanosText, 10, 64)
|
||||
if err != nil {
|
||||
return t, err
|
||||
}
|
||||
}
|
||||
|
||||
t = time.Unix(sec, nanos)
|
||||
return t, nil
|
||||
}
|
||||
|
||||
// Decodes a value from JSON, coercing it to a string value as necessary
|
||||
func decodeAsString(dec *json.Decoder) (s string, err error) {
|
||||
var v interface{}
|
||||
if err = dec.Decode(&v); err != nil {
|
||||
return s, err
|
||||
}
|
||||
var ok bool
|
||||
if s, ok = v.(string); ok {
|
||||
return s, err
|
||||
}
|
||||
s = fmt.Sprint(v)
|
||||
return s, err
|
||||
}
|
||||
@@ -1,23 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package human
|
||||
|
||||
import (
|
||||
"bufio"
|
||||
"io"
|
||||
)
|
||||
|
||||
type LogWriter interface {
|
||||
Write(e LogEntry)
|
||||
}
|
||||
|
||||
// Read JSON logs from input and write formatted logs to the output
|
||||
func ProcessLogs(reader io.Reader, writer LogWriter) {
|
||||
scanner := bufio.NewScanner(reader)
|
||||
for scanner.Scan() {
|
||||
s := scanner.Text()
|
||||
e := ParseLogMessage(s)
|
||||
writer.Write(e)
|
||||
}
|
||||
}
|
||||
@@ -1,23 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package human
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"io"
|
||||
)
|
||||
|
||||
type SimpleWriter struct {
|
||||
out io.Writer
|
||||
}
|
||||
|
||||
func (w *SimpleWriter) Write(e LogEntry) {
|
||||
fmt.Fprintln(w.out, e)
|
||||
}
|
||||
|
||||
func NewSimpleWriter(out io.Writer) *SimpleWriter {
|
||||
w := new(SimpleWriter)
|
||||
w.out = out
|
||||
return w
|
||||
}
|
||||
@@ -3,49 +3,56 @@
|
||||
|
||||
package mlog
|
||||
|
||||
// Standard levels
|
||||
import "github.com/mattermost/logr/v2"
|
||||
|
||||
// Standard levels.
|
||||
var (
|
||||
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"}
|
||||
LvlPanic = logr.Panic // ID = 0
|
||||
LvlFatal = logr.Fatal // ID = 1
|
||||
LvlError = logr.Error // ID = 2
|
||||
LvlWarn = logr.Warn // ID = 3
|
||||
LvlInfo = logr.Info // ID = 4
|
||||
LvlDebug = logr.Debug // ID = 5
|
||||
LvlTrace = logr.Trace // ID = 6
|
||||
StdAll = []Level{LvlPanic, LvlFatal, LvlError, LvlWarn, LvlInfo, LvlDebug, LvlTrace}
|
||||
// non-standard "critical" level
|
||||
LvlCritical = Level{ID: 7, Name: "critical"}
|
||||
// used by redirected standard logger
|
||||
LvlStdLog = LogLevel{ID: 10, Name: "stdlog"}
|
||||
LvlStdLog = Level{ID: 10, Name: "stdlog"}
|
||||
// used only by the logger
|
||||
LvlLogError = LogLevel{ID: 11, Name: "logerror", Stacktrace: true}
|
||||
LvlLogError = Level{ID: 11, Name: "logerror", Stacktrace: true}
|
||||
)
|
||||
|
||||
// Register custom (discrete) levels here.
|
||||
// !!!!! ID's must not exceed 32,768 !!!!!!
|
||||
// !!!!! Custom ID's must be between 20 and 32,768 !!!!!!
|
||||
var (
|
||||
// used by the audit system
|
||||
LvlAuditAPI = LogLevel{ID: 100, Name: "audit-api"}
|
||||
LvlAuditContent = LogLevel{ID: 101, Name: "audit-content"}
|
||||
LvlAuditPerms = LogLevel{ID: 102, Name: "audit-permissions"}
|
||||
LvlAuditCLI = LogLevel{ID: 103, Name: "audit-cli"}
|
||||
LvlAuditAPI = Level{ID: 100, Name: "audit-api"}
|
||||
LvlAuditContent = Level{ID: 101, Name: "audit-content"}
|
||||
LvlAuditPerms = Level{ID: 102, Name: "audit-permissions"}
|
||||
LvlAuditCLI = Level{ID: 103, Name: "audit-cli"}
|
||||
|
||||
// used by the TCP log target
|
||||
LvlTCPLogTarget = LogLevel{ID: 120, Name: "TcpLogTarget"}
|
||||
LvlTCPLogTarget = Level{ID: 120, Name: "TcpLogTarget"}
|
||||
|
||||
// used by Remote Cluster Service
|
||||
LvlRemoteClusterServiceDebug = LogLevel{ID: 130, Name: "RemoteClusterServiceDebug"}
|
||||
LvlRemoteClusterServiceError = LogLevel{ID: 131, Name: "RemoteClusterServiceError"}
|
||||
LvlRemoteClusterServiceWarn = LogLevel{ID: 132, Name: "RemoteClusterServiceWarn"}
|
||||
LvlRemoteClusterServiceDebug = Level{ID: 130, Name: "RemoteClusterServiceDebug"}
|
||||
LvlRemoteClusterServiceError = Level{ID: 131, Name: "RemoteClusterServiceError"}
|
||||
LvlRemoteClusterServiceWarn = Level{ID: 132, Name: "RemoteClusterServiceWarn"}
|
||||
|
||||
// used by Shared Channel Sync Service
|
||||
LvlSharedChannelServiceDebug = LogLevel{ID: 200, Name: "SharedChannelServiceDebug"}
|
||||
LvlSharedChannelServiceError = LogLevel{ID: 201, Name: "SharedChannelServiceError"}
|
||||
LvlSharedChannelServiceWarn = LogLevel{ID: 202, Name: "SharedChannelServiceWarn"}
|
||||
LvlSharedChannelServiceMessagesInbound = LogLevel{ID: 203, Name: "SharedChannelServiceMsgInbound"}
|
||||
LvlSharedChannelServiceMessagesOutbound = LogLevel{ID: 204, Name: "SharedChannelServiceMsgOutbound"}
|
||||
LvlSharedChannelServiceDebug = Level{ID: 200, Name: "SharedChannelServiceDebug"}
|
||||
LvlSharedChannelServiceError = Level{ID: 201, Name: "SharedChannelServiceError"}
|
||||
LvlSharedChannelServiceWarn = Level{ID: 202, Name: "SharedChannelServiceWarn"}
|
||||
LvlSharedChannelServiceMessagesInbound = Level{ID: 203, Name: "SharedChannelServiceMsgInbound"}
|
||||
LvlSharedChannelServiceMessagesOutbound = Level{ID: 204, Name: "SharedChannelServiceMsgOutbound"}
|
||||
|
||||
// add more here ...
|
||||
// Focalboard
|
||||
LvlFBTelemetry = Level{ID: 9000, Name: "telemetry"}
|
||||
LvlFBMetrics = Level{ID: 9001, Name: "metrics"}
|
||||
)
|
||||
|
||||
// Combinations for LogM (log multi)
|
||||
// Combinations for LogM (log multi).
|
||||
var (
|
||||
MLvlAuditAll = []LogLevel{LvlAuditAPI, LvlAuditContent, LvlAuditPerms, LvlAuditCLI}
|
||||
MLvlAuditAll = []Level{LvlAuditAPI, LvlAuditContent, LvlAuditPerms, LvlAuditCLI}
|
||||
)
|
||||
|
||||
@@ -1,361 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"os"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/mattermost/logr"
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
"gopkg.in/natefinch/lumberjack.v2"
|
||||
)
|
||||
|
||||
const (
|
||||
// Very verbose messages for debugging specific issues
|
||||
LevelDebug = "debug"
|
||||
// Default log level, informational
|
||||
LevelInfo = "info"
|
||||
// Warnings are messages about possible issues
|
||||
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
|
||||
type Field = zapcore.Field
|
||||
|
||||
var Int64 = zap.Int64
|
||||
var Int32 = zap.Int32
|
||||
var Int = zap.Int
|
||||
var Uint32 = zap.Uint32
|
||||
var String = zap.String
|
||||
var Any = zap.Any
|
||||
var Err = zap.Error
|
||||
var NamedErr = zap.NamedError
|
||||
var Bool = zap.Bool
|
||||
var Duration = zap.Duration
|
||||
|
||||
type LoggerIFace interface {
|
||||
IsLevelEnabled(LogLevel) bool
|
||||
Debug(string, ...Field)
|
||||
Info(string, ...Field)
|
||||
Warn(string, ...Field)
|
||||
Error(string, ...Field)
|
||||
Critical(string, ...Field)
|
||||
Log(LogLevel, string, ...Field)
|
||||
LogM([]LogLevel, string, ...Field)
|
||||
}
|
||||
|
||||
type TargetInfo logr.TargetInfo
|
||||
|
||||
type LoggerConfiguration struct {
|
||||
EnableConsole bool
|
||||
ConsoleJson bool
|
||||
EnableColor bool
|
||||
ConsoleLevel string
|
||||
EnableFile bool
|
||||
FileJson bool
|
||||
FileLevel string
|
||||
FileLocation string
|
||||
}
|
||||
|
||||
type Logger struct {
|
||||
zap *zap.Logger
|
||||
consoleLevel zap.AtomicLevel
|
||||
fileLevel zap.AtomicLevel
|
||||
logrLogger *logr.Logger
|
||||
mutex *sync.RWMutex
|
||||
}
|
||||
|
||||
func getZapLevel(level string) zapcore.Level {
|
||||
switch level {
|
||||
case LevelInfo:
|
||||
return zapcore.InfoLevel
|
||||
case LevelWarn:
|
||||
return zapcore.WarnLevel
|
||||
case LevelDebug:
|
||||
return zapcore.DebugLevel
|
||||
case LevelError:
|
||||
return zapcore.ErrorLevel
|
||||
default:
|
||||
return zapcore.InfoLevel
|
||||
}
|
||||
}
|
||||
|
||||
func makeEncoder(json, color bool) zapcore.Encoder {
|
||||
encoderConfig := zap.NewProductionEncoderConfig()
|
||||
if json {
|
||||
return zapcore.NewJSONEncoder(encoderConfig)
|
||||
}
|
||||
|
||||
if color {
|
||||
encoderConfig.EncodeLevel = zapcore.CapitalColorLevelEncoder
|
||||
}
|
||||
encoderConfig.EncodeTime = zapcore.ISO8601TimeEncoder
|
||||
return zapcore.NewConsoleEncoder(encoderConfig)
|
||||
}
|
||||
|
||||
func NewLogger(config *LoggerConfiguration) *Logger {
|
||||
cores := []zapcore.Core{}
|
||||
logger := &Logger{
|
||||
consoleLevel: zap.NewAtomicLevelAt(getZapLevel(config.ConsoleLevel)),
|
||||
fileLevel: zap.NewAtomicLevelAt(getZapLevel(config.FileLevel)),
|
||||
logrLogger: newLogr(),
|
||||
mutex: &sync.RWMutex{},
|
||||
}
|
||||
|
||||
if config.EnableConsole {
|
||||
writer := zapcore.Lock(os.Stderr)
|
||||
core := zapcore.NewCore(makeEncoder(config.ConsoleJson, config.EnableColor), writer, logger.consoleLevel)
|
||||
cores = append(cores, core)
|
||||
}
|
||||
|
||||
if config.EnableFile {
|
||||
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, false), writer, logger.fileLevel)
|
||||
cores = append(cores, core)
|
||||
}
|
||||
}
|
||||
|
||||
combinedCore := zapcore.NewTee(cores...)
|
||||
|
||||
logger.zap = zap.New(combinedCore,
|
||||
zap.AddCaller(),
|
||||
)
|
||||
return logger
|
||||
}
|
||||
|
||||
func (l *Logger) ChangeLevels(config *LoggerConfiguration) {
|
||||
l.consoleLevel.SetLevel(getZapLevel(config.ConsoleLevel))
|
||||
l.fileLevel.SetLevel(getZapLevel(config.FileLevel))
|
||||
}
|
||||
|
||||
func (l *Logger) SetConsoleLevel(level string) {
|
||||
l.consoleLevel.SetLevel(getZapLevel(level))
|
||||
}
|
||||
|
||||
func (l *Logger) With(fields ...Field) *Logger {
|
||||
newLogger := *l
|
||||
newLogger.zap = newLogger.zap.With(fields...)
|
||||
if newLogger.getLogger() != nil {
|
||||
ll := newLogger.getLogger().WithFields(zapToLogr(fields))
|
||||
newLogger.logrLogger = &ll
|
||||
}
|
||||
return &newLogger
|
||||
}
|
||||
|
||||
func (l *Logger) StdLog(fields ...Field) *log.Logger {
|
||||
return zap.NewStdLog(l.With(fields...).zap.WithOptions(getStdLogOption()))
|
||||
}
|
||||
|
||||
// StdLogAt returns *log.Logger which writes to supplied zap logger at required level.
|
||||
func (l *Logger) StdLogAt(level string, fields ...Field) (*log.Logger, error) {
|
||||
return zap.NewStdLogAt(l.With(fields...).zap.WithOptions(getStdLogOption()), getZapLevel(level))
|
||||
}
|
||||
|
||||
// StdLogWriter returns a writer that can be hooked up to the output of a golang standard logger
|
||||
// anything written will be interpreted as log entries accordingly
|
||||
func (l *Logger) StdLogWriter() io.Writer {
|
||||
newLogger := *l
|
||||
newLogger.zap = newLogger.zap.WithOptions(zap.AddCallerSkip(4), getStdLogOption())
|
||||
f := newLogger.Info
|
||||
return &loggerWriter{f}
|
||||
}
|
||||
|
||||
func (l *Logger) WithCallerSkip(skip int) *Logger {
|
||||
newLogger := *l
|
||||
newLogger.zap = newLogger.zap.WithOptions(zap.AddCallerSkip(skip))
|
||||
return &newLogger
|
||||
}
|
||||
|
||||
// Made for the plugin interface, wraps mlog in a simpler interface
|
||||
// at the cost of performance
|
||||
func (l *Logger) Sugar() *SugarLogger {
|
||||
return &SugarLogger{
|
||||
wrappedLogger: l,
|
||||
zapSugar: l.zap.Sugar(),
|
||||
}
|
||||
}
|
||||
|
||||
func (l *Logger) IsLevelEnabled(level LogLevel) bool {
|
||||
return isLevelEnabled(l.getLogger(), logr.Level(level))
|
||||
}
|
||||
|
||||
func (l *Logger) Debug(message string, fields ...Field) {
|
||||
l.zap.Debug(message, fields...)
|
||||
if isLevelEnabled(l.getLogger(), logr.Debug) {
|
||||
l.getLogger().WithFields(zapToLogr(fields)).Debug(message)
|
||||
}
|
||||
}
|
||||
|
||||
func (l *Logger) Info(message string, fields ...Field) {
|
||||
l.zap.Info(message, fields...)
|
||||
if isLevelEnabled(l.getLogger(), logr.Info) {
|
||||
l.getLogger().WithFields(zapToLogr(fields)).Info(message)
|
||||
}
|
||||
}
|
||||
|
||||
func (l *Logger) Warn(message string, fields ...Field) {
|
||||
l.zap.Warn(message, fields...)
|
||||
if isLevelEnabled(l.getLogger(), logr.Warn) {
|
||||
l.getLogger().WithFields(zapToLogr(fields)).Warn(message)
|
||||
}
|
||||
}
|
||||
|
||||
func (l *Logger) Error(message string, fields ...Field) {
|
||||
l.zap.Error(message, fields...)
|
||||
if isLevelEnabled(l.getLogger(), logr.Error) {
|
||||
l.getLogger().WithFields(zapToLogr(fields)).Error(message)
|
||||
}
|
||||
}
|
||||
|
||||
func (l *Logger) Critical(message string, fields ...Field) {
|
||||
l.zap.Error(message, fields...)
|
||||
if isLevelEnabled(l.getLogger(), logr.Error) {
|
||||
l.getLogger().WithFields(zapToLogr(fields)).Error(message)
|
||||
}
|
||||
}
|
||||
|
||||
func (l *Logger) Log(level LogLevel, message string, fields ...Field) {
|
||||
l.getLogger().WithFields(zapToLogr(fields)).Log(logr.Level(level), message)
|
||||
}
|
||||
|
||||
func (l *Logger) LogM(levels []LogLevel, message string, fields ...Field) {
|
||||
var logger *logr.Logger
|
||||
for _, lvl := range levels {
|
||||
if isLevelEnabled(l.getLogger(), logr.Level(lvl)) {
|
||||
// don't create logger with fields unless at least one level is active.
|
||||
if logger == nil {
|
||||
l := l.getLogger().WithFields(zapToLogr(fields))
|
||||
logger = &l
|
||||
}
|
||||
logger.Log(logr.Level(lvl), message)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (l *Logger) Flush(cxt context.Context) error {
|
||||
return l.getLogger().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.getLogger().Logr().ShutdownWithTimeout(cxt)
|
||||
l.setLogger(newLogr())
|
||||
return err
|
||||
}
|
||||
|
||||
// ConfigAdvancedLoggingConfig (re)configures advanced logging based on the
|
||||
// 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 err := l.ShutdownAdvancedLogging(context.Background()); err != nil {
|
||||
Error("error shutting down previous logger", Err(err))
|
||||
}
|
||||
|
||||
err := logrAddTargets(l.getLogger(), targets)
|
||||
return err
|
||||
}
|
||||
|
||||
// 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(targets ...logr.Target) error {
|
||||
return l.getLogger().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.getLogger().Logr().RemoveTargets(ctx, fc)
|
||||
}
|
||||
|
||||
// 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.getLogger().Logr().SetMetricsCollector(collector)
|
||||
}
|
||||
|
||||
// getLogger is a concurrent safe getter of the logr logger
|
||||
func (l *Logger) getLogger() *logr.Logger {
|
||||
defer l.mutex.RUnlock()
|
||||
l.mutex.RLock()
|
||||
return l.logrLogger
|
||||
}
|
||||
|
||||
// setLogger is a concurrent safe setter of the logr logger
|
||||
func (l *Logger) setLogger(logger *logr.Logger) {
|
||||
defer l.mutex.Unlock()
|
||||
l.mutex.Lock()
|
||||
l.logrLogger = logger
|
||||
}
|
||||
|
||||
// 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)
|
||||
}
|
||||
@@ -1,51 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"github.com/mattermost/mattermost-server/v6/shared/mlog"
|
||||
)
|
||||
|
||||
// Test race condition when shutting down advanced logging. This test must run with the -race flag in order to verify
|
||||
// that there is no race.
|
||||
func TestLogger_ShutdownAdvancedLoggingRace(t *testing.T) {
|
||||
logger := mlog.NewLogger(&mlog.LoggerConfiguration{
|
||||
EnableConsole: true,
|
||||
ConsoleJson: true,
|
||||
EnableFile: false,
|
||||
FileLevel: mlog.LevelInfo,
|
||||
})
|
||||
started := make(chan bool)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
var wg sync.WaitGroup
|
||||
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
started <- true
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return
|
||||
default:
|
||||
logger.Debug("testing...")
|
||||
}
|
||||
}
|
||||
}()
|
||||
|
||||
<-started
|
||||
|
||||
err := logger.ShutdownAdvancedLogging(ctx)
|
||||
require.NoError(t, err)
|
||||
|
||||
cancel()
|
||||
wg.Wait()
|
||||
}
|
||||
@@ -1,244 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"os"
|
||||
|
||||
"github.com/hashicorp/go-multierror"
|
||||
"github.com/mattermost/logr"
|
||||
logrFmt "github.com/mattermost/logr/format"
|
||||
"github.com/mattermost/logr/target"
|
||||
"go.uber.org/zap/zapcore"
|
||||
)
|
||||
|
||||
const (
|
||||
DefaultMaxTargetQueue = 1000
|
||||
DefaultSysLogPort = 514
|
||||
)
|
||||
|
||||
type LogLevel struct {
|
||||
ID logr.LevelID
|
||||
Name string
|
||||
Stacktrace bool
|
||||
}
|
||||
|
||||
type LogTarget struct {
|
||||
Type string // one of "console", "file", "tcp", "syslog", "none".
|
||||
Format string // one of "json", "plain"
|
||||
Levels []LogLevel
|
||||
Options json.RawMessage
|
||||
MaxQueueSize int
|
||||
}
|
||||
|
||||
type LogTargetCfg map[string]*LogTarget
|
||||
type LogrCleanup func() error
|
||||
|
||||
func newLogr() *logr.Logger {
|
||||
lgr := &logr.Logr{}
|
||||
lgr.OnExit = func(int) {}
|
||||
lgr.OnPanic = func(interface{}) {}
|
||||
lgr.OnLoggerError = onLoggerError
|
||||
lgr.OnQueueFull = onQueueFull
|
||||
lgr.OnTargetQueueFull = onTargetQueueFull
|
||||
|
||||
logger := lgr.NewLogger()
|
||||
return &logger
|
||||
}
|
||||
|
||||
func logrAddTargets(logger *logr.Logger, targets LogTargetCfg) error {
|
||||
lgr := logger.Logr()
|
||||
var errs error
|
||||
for name, t := range targets {
|
||||
target, err := NewLogrTarget(name, t)
|
||||
if err != nil {
|
||||
errs = multierror.Append(err)
|
||||
continue
|
||||
}
|
||||
if target != nil {
|
||||
target.SetName(name)
|
||||
lgr.AddTarget(target)
|
||||
}
|
||||
}
|
||||
return errs
|
||||
}
|
||||
|
||||
// NewLogrTarget creates a `logr.Target` based on a target config.
|
||||
// Can be used when parsing custom config files, or when programmatically adding
|
||||
// built-in targets. Use `mlog.AddTarget` to add custom targets.
|
||||
func NewLogrTarget(name string, t *LogTarget) (logr.Target, error) {
|
||||
formatter, err := newFormatter(name, t.Format)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
filter := newFilter(t.Levels)
|
||||
|
||||
if t.MaxQueueSize == 0 {
|
||||
t.MaxQueueSize = DefaultMaxTargetQueue
|
||||
}
|
||||
|
||||
switch t.Type {
|
||||
case "console":
|
||||
return newConsoleTarget(name, t, filter, formatter)
|
||||
case "file":
|
||||
return newFileTarget(name, t, filter, formatter)
|
||||
case "syslog":
|
||||
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)
|
||||
}
|
||||
|
||||
func newFilter(levels []LogLevel) logr.Filter {
|
||||
filter := &logr.CustomFilter{}
|
||||
for _, lvl := range levels {
|
||||
filter.Add(logr.Level(lvl))
|
||||
}
|
||||
return filter
|
||||
}
|
||||
|
||||
func newFormatter(name string, format string) (logr.Formatter, error) {
|
||||
switch format {
|
||||
case "json", "":
|
||||
return &logrFmt.JSON{}, nil
|
||||
case "plain":
|
||||
return &logrFmt.Plain{Delim: " | "}, nil
|
||||
default:
|
||||
return nil, fmt.Errorf("invalid format '%s' for target %s", format, name)
|
||||
}
|
||||
}
|
||||
|
||||
func newConsoleTarget(name string, t *LogTarget, filter logr.Filter, formatter logr.Formatter) (logr.Target, error) {
|
||||
type consoleOptions struct {
|
||||
Out string `json:"Out"`
|
||||
}
|
||||
options := &consoleOptions{}
|
||||
if err := json.Unmarshal(t.Options, options); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
var w io.Writer
|
||||
switch options.Out {
|
||||
case "stdout", "":
|
||||
w = os.Stdout
|
||||
case "stderr":
|
||||
w = os.Stderr
|
||||
default:
|
||||
return nil, fmt.Errorf("invalid out '%s' for target %s", options.Out, name)
|
||||
}
|
||||
|
||||
newTarget := target.NewWriterTarget(filter, formatter, w, t.MaxQueueSize)
|
||||
return newTarget, nil
|
||||
}
|
||||
|
||||
func newFileTarget(name string, t *LogTarget, filter logr.Filter, formatter logr.Formatter) (logr.Target, error) {
|
||||
type fileOptions struct {
|
||||
Filename string `json:"Filename"`
|
||||
MaxSize int `json:"MaxSizeMB"`
|
||||
MaxAge int `json:"MaxAgeDays"`
|
||||
MaxBackups int `json:"MaxBackups"`
|
||||
Compress bool `json:"Compress"`
|
||||
}
|
||||
options := &fileOptions{}
|
||||
if err := json.Unmarshal(t.Options, options); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return newFileTargetWithOpts(name, t, target.FileOptions(*options), filter, formatter)
|
||||
}
|
||||
|
||||
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(opts.Filename); err != nil {
|
||||
return nil, fmt.Errorf("error writing to 'Filename' for target %s: %w", name, err)
|
||||
}
|
||||
|
||||
newTarget := target.NewFileTarget(filter, formatter, opts, t.MaxQueueSize)
|
||||
return newTarget, nil
|
||||
}
|
||||
|
||||
func newSyslogTarget(name string, t *LogTarget, filter logr.Filter, formatter logr.Formatter) (logr.Target, error) {
|
||||
options := &SyslogParams{}
|
||||
if err := json.Unmarshal(t.Options, options); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if options.IP == "" {
|
||||
return nil, fmt.Errorf("missing 'IP' option for target %s", name)
|
||||
}
|
||||
if options.Port == 0 {
|
||||
options.Port = DefaultSysLogPort
|
||||
}
|
||||
return NewSyslogTarget(filter, formatter, options, t.MaxQueueSize)
|
||||
}
|
||||
|
||||
func newTCPTarget(name string, t *LogTarget, filter logr.Filter, formatter logr.Formatter) (logr.Target, error) {
|
||||
options := &TCPParams{}
|
||||
if err := json.Unmarshal(t.Options, options); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if options.IP == "" {
|
||||
return nil, fmt.Errorf("missing 'IP' option for target %s", name)
|
||||
}
|
||||
if options.Port == 0 {
|
||||
return nil, fmt.Errorf("missing 'Port' option for target %s", name)
|
||||
}
|
||||
return NewTCPTarget(filter, formatter, options, t.MaxQueueSize)
|
||||
}
|
||||
|
||||
func checkFileWritable(filename string) error {
|
||||
// try opening/creating the file for writing
|
||||
file, err := os.OpenFile(filename, os.O_RDWR|os.O_APPEND|os.O_CREATE, 0600)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
file.Close()
|
||||
return nil
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
|
||||
// zapToLogr converts Zap fields to Logr fields.
|
||||
// This will not be needed once Logr is used for all logging.
|
||||
func zapToLogr(zapFields []Field) logr.Fields {
|
||||
encoder := zapcore.NewMapObjectEncoder()
|
||||
for _, zapField := range zapFields {
|
||||
zapField.AddTo(encoder)
|
||||
}
|
||||
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
|
||||
}
|
||||
407
shared/mlog/mlog.go
Обычный файл
407
shared/mlog/mlog.go
Обычный файл
@@ -0,0 +1,407 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
// Package mlog provides a simple wrapper around Logr.
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"log"
|
||||
"os"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"github.com/mattermost/logr/v2"
|
||||
logrcfg "github.com/mattermost/logr/v2/config"
|
||||
)
|
||||
|
||||
const (
|
||||
ShutdownTimeout = time.Second * 15
|
||||
FlushTimeout = time.Second * 15
|
||||
DefaultMaxQueueSize = 1000
|
||||
DefaultMetricsUpdateFreqMillis = 15000
|
||||
)
|
||||
|
||||
type LoggerIFace interface {
|
||||
IsLevelEnabled(Level) bool
|
||||
Debug(string, ...Field)
|
||||
Info(string, ...Field)
|
||||
Warn(string, ...Field)
|
||||
Error(string, ...Field)
|
||||
Critical(string, ...Field)
|
||||
Log(Level, string, ...Field)
|
||||
LogM([]Level, string, ...Field)
|
||||
}
|
||||
|
||||
// Type and function aliases from Logr to limit the spread of dependencies.
|
||||
type Field = logr.Field
|
||||
type Level = logr.Level
|
||||
type Option = logr.Option
|
||||
type Target = logr.Target
|
||||
type TargetInfo = logr.TargetInfo
|
||||
type LogRec = logr.LogRec
|
||||
type LogCloner = logr.LogCloner
|
||||
type MetricsCollector = logr.MetricsCollector
|
||||
type TargetCfg = logrcfg.TargetCfg
|
||||
type Sugar = logr.Sugar
|
||||
|
||||
// LoggerConfiguration is a map of LogTarget configurations.
|
||||
type LoggerConfiguration map[string]TargetCfg
|
||||
|
||||
func (lc LoggerConfiguration) Append(cfg LoggerConfiguration) {
|
||||
for k, v := range cfg {
|
||||
lc[k] = v
|
||||
}
|
||||
}
|
||||
|
||||
func (lc LoggerConfiguration) toTargetCfg() map[string]logrcfg.TargetCfg {
|
||||
tcfg := make(map[string]logrcfg.TargetCfg)
|
||||
for k, v := range lc {
|
||||
tcfg[k] = v
|
||||
}
|
||||
return tcfg
|
||||
}
|
||||
|
||||
// Any picks the best supported field type based on type of val.
|
||||
// For best performance when passing a struct (or struct pointer),
|
||||
// implement `logr.LogWriter` on the struct, otherwise reflection
|
||||
// will be used to generate a string representation.
|
||||
var Any = logr.Any
|
||||
|
||||
// Int64 constructs a field containing a key and Int64 value.
|
||||
var Int64 = logr.Int64
|
||||
|
||||
// Int32 constructs a field containing a key and Int32 value.
|
||||
var Int32 = logr.Int32
|
||||
|
||||
// Int constructs a field containing a key and Int value.
|
||||
var Int = logr.Int
|
||||
|
||||
// Uint64 constructs a field containing a key and Uint64 value.
|
||||
var Uint64 = logr.Uint64
|
||||
|
||||
// Uint32 constructs a field containing a key and Uint32 value.
|
||||
var Uint32 = logr.Uint32
|
||||
|
||||
// Uint constructs a field containing a key and Uint value.
|
||||
var Uint = logr.Uint
|
||||
|
||||
// Float64 constructs a field containing a key and Float64 value.
|
||||
var Float64 = logr.Float64
|
||||
|
||||
// Float32 constructs a field containing a key and Float32 value.
|
||||
var Float32 = logr.Float32
|
||||
|
||||
// String constructs a field containing a key and String value.
|
||||
var String = logr.String
|
||||
|
||||
// Stringer constructs a field containing a key and a fmt.Stringer value.
|
||||
// The fmt.Stringer's `String` method is called lazily.
|
||||
var Stringer = logr.Stringer
|
||||
|
||||
// Err constructs a field containing a default key ("error") and error value.
|
||||
var Err = logr.Err
|
||||
|
||||
// NamedErr constructs a field containing a key and error value.
|
||||
var NamedErr = logr.NamedErr
|
||||
|
||||
// Bool constructs a field containing a key and bool value.
|
||||
var Bool = logr.Bool
|
||||
|
||||
// Time constructs a field containing a key and time.Time value.
|
||||
var Time = logr.Time
|
||||
|
||||
// Duration constructs a field containing a key and time.Duration value.
|
||||
var Duration = logr.Duration
|
||||
|
||||
// Millis constructs a field containing a key and timestamp value.
|
||||
// The timestamp is expected to be milliseconds since Jan 1, 1970 UTC.
|
||||
var Millis = logr.Millis
|
||||
|
||||
// Array constructs a field containing a key and array value.
|
||||
var Array = logr.Array
|
||||
|
||||
// Map constructs a field containing a key and map value.
|
||||
var Map = logr.Map
|
||||
|
||||
// Logger provides a thin wrapper around a Logr instance. This is a struct instead of an interface
|
||||
// so that there are no allocations on the heap each interface method invocation. Normally not
|
||||
// something to be concerned about, but logging calls for disabled levels should have as little CPU
|
||||
// and memory impact as possible. Most of these wrapper calls will be inlined as well.
|
||||
type Logger struct {
|
||||
log *logr.Logger
|
||||
lockConfig *int32
|
||||
}
|
||||
|
||||
// NewLogger creates a new Logger instance which can be configured via `(*Logger).Configure`.
|
||||
// Some options with invalid values can cause an error to be returned, however `NewLogger()`
|
||||
// using just defaults never errors.
|
||||
func NewLogger(options ...Option) (*Logger, error) {
|
||||
options = append(options, logr.StackFilter(logr.GetPackageName("NewLogger")))
|
||||
|
||||
lgr, err := logr.New(options...)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
log := lgr.NewLogger()
|
||||
var lockConfig int32
|
||||
|
||||
return &Logger{
|
||||
log: &log,
|
||||
lockConfig: &lockConfig,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Configure provides a new configuration for this logger.
|
||||
// Zero or more sources of config can be provided:
|
||||
// cfgFile - path to file containing JSON
|
||||
// cfgEscaped - JSON string probably from ENV var
|
||||
//
|
||||
// For each case JSON containing log targets is provided. Target name collisions are resolved
|
||||
// using the following precedence:
|
||||
// cfgFile > cfgEscaped
|
||||
func (l *Logger) Configure(cfgFile string, cfgEscaped string) error {
|
||||
if atomic.LoadInt32(l.lockConfig) != 0 {
|
||||
return ErrConfigurationLock
|
||||
}
|
||||
|
||||
cfgMap := make(LoggerConfiguration)
|
||||
|
||||
// Add config from file
|
||||
if cfgFile != "" {
|
||||
b, err := ioutil.ReadFile(cfgFile)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error reading logger config file %s: %w", cfgFile, err)
|
||||
}
|
||||
|
||||
var mapCfgFile LoggerConfiguration
|
||||
if err := json.Unmarshal(b, &mapCfgFile); err != nil {
|
||||
return fmt.Errorf("error decoding logger config file %s: %w", cfgFile, err)
|
||||
}
|
||||
cfgMap.Append(mapCfgFile)
|
||||
}
|
||||
|
||||
// Add config from escaped json string
|
||||
if cfgEscaped != "" {
|
||||
var mapCfgEscaped LoggerConfiguration
|
||||
if err := json.Unmarshal([]byte(cfgEscaped), &mapCfgEscaped); err != nil {
|
||||
return fmt.Errorf("error decoding logger config as escaped json: %w", err)
|
||||
}
|
||||
cfgMap.Append(mapCfgEscaped)
|
||||
}
|
||||
|
||||
if len(cfgMap) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
return logrcfg.ConfigureTargets(l.log.Logr(), cfgMap.toTargetCfg(), nil)
|
||||
}
|
||||
|
||||
// ConfigureTargets provides a new configuration for this logger via a `LoggerConfig` map.
|
||||
// Typically `mlog.Configure` is used instead which accepts JSON formatted configuration.
|
||||
func (l *Logger) ConfigureTargets(cfg LoggerConfiguration) error {
|
||||
if atomic.LoadInt32(l.lockConfig) != 0 {
|
||||
return ErrConfigurationLock
|
||||
}
|
||||
return logrcfg.ConfigureTargets(l.log.Logr(), cfg.toTargetCfg(), nil)
|
||||
}
|
||||
|
||||
// LockConfiguration disallows further configuration changes until `UnlockConfiguration`
|
||||
// is called. The previous locked stated is returned.
|
||||
func (l *Logger) LockConfiguration() bool {
|
||||
old := atomic.SwapInt32(l.lockConfig, 1)
|
||||
return old != 0
|
||||
}
|
||||
|
||||
// UnlockConfiguration allows configuration changes. The previous locked stated is returned.
|
||||
func (l *Logger) UnlockConfiguration() bool {
|
||||
old := atomic.SwapInt32(l.lockConfig, 0)
|
||||
return old != 0
|
||||
}
|
||||
|
||||
// IsConfigurationLocked returns the current state of the configuration lock.
|
||||
func (l *Logger) IsConfigurationLocked() bool {
|
||||
return atomic.LoadInt32(l.lockConfig) != 0
|
||||
}
|
||||
|
||||
// With creates a new Logger with the specified fields. This is a light-weight
|
||||
// operation and can be called on demand.
|
||||
func (l *Logger) With(fields ...Field) *Logger {
|
||||
logWith := l.log.With(fields...)
|
||||
return &Logger{
|
||||
log: &logWith,
|
||||
lockConfig: l.lockConfig,
|
||||
}
|
||||
}
|
||||
|
||||
// IsLevelEnabled returns true only if at least one log target is
|
||||
// configured to emit the specified log level. Use this check when
|
||||
// gathering the log info may be expensive.
|
||||
//
|
||||
// Note, transformations and serializations done via fields are already
|
||||
// lazily evaluated and don't require this check beforehand.
|
||||
func (l *Logger) IsLevelEnabled(level Level) bool {
|
||||
return l.log.IsLevelEnabled(level)
|
||||
}
|
||||
|
||||
// Log emits the log record for any targets configured for the specified level.
|
||||
func (l *Logger) Log(level Level, msg string, fields ...Field) {
|
||||
l.log.Log(level, msg, fields...)
|
||||
}
|
||||
|
||||
// LogM emits the log record for any targets configured for the specified levels.
|
||||
// Equivalent to calling `Log` once for each level.
|
||||
func (l *Logger) LogM(levels []Level, msg string, fields ...Field) {
|
||||
l.log.LogM(levels, msg, fields...)
|
||||
}
|
||||
|
||||
// Convenience method equivalent to calling `Log` with the `Trace` level.
|
||||
func (l *Logger) Trace(msg string, fields ...Field) {
|
||||
l.log.Trace(msg, fields...)
|
||||
}
|
||||
|
||||
// Convenience method equivalent to calling `Log` with the `Debug` level.
|
||||
func (l *Logger) Debug(msg string, fields ...Field) {
|
||||
l.log.Debug(msg, fields...)
|
||||
}
|
||||
|
||||
// Convenience method equivalent to calling `Log` with the `Info` level.
|
||||
func (l *Logger) Info(msg string, fields ...Field) {
|
||||
l.log.Info(msg, fields...)
|
||||
}
|
||||
|
||||
// Convenience method equivalent to calling `Log` with the `Warn` level.
|
||||
func (l *Logger) Warn(msg string, fields ...Field) {
|
||||
l.log.Warn(msg, fields...)
|
||||
}
|
||||
|
||||
// Convenience method equivalent to calling `Log` with the `Error` level.
|
||||
func (l *Logger) Error(msg string, fields ...Field) {
|
||||
l.log.Error(msg, fields...)
|
||||
}
|
||||
|
||||
// Convenience method equivalent to calling `Log` with the `Critical` level.
|
||||
func (l *Logger) Critical(msg string, fields ...Field) {
|
||||
l.log.Log(LvlCritical, msg, fields...)
|
||||
}
|
||||
|
||||
// Convenience method equivalent to calling `Log` with the `Fatal` level,
|
||||
// followed by `os.Exit(1)`.
|
||||
func (l *Logger) Fatal(msg string, fields ...Field) {
|
||||
l.log.Log(logr.Fatal, msg, fields...)
|
||||
_ = l.Shutdown()
|
||||
os.Exit(1)
|
||||
}
|
||||
|
||||
// HasTargets returns true if at least one log target has been added.
|
||||
func (l *Logger) HasTargets() bool {
|
||||
return l.log.Logr().HasTargets()
|
||||
}
|
||||
|
||||
// StdLogger creates a standard logger backed by this logger.
|
||||
// All log records are output with the specified level.
|
||||
func (l *Logger) StdLogger(level Level) *log.Logger {
|
||||
return l.log.StdLogger(level)
|
||||
}
|
||||
|
||||
// StdLogWriter returns a writer that can be hooked up to the output of a golang standard logger
|
||||
// anything written will be interpreted as log entries and passed to this logger.
|
||||
func (l *Logger) StdLogWriter() io.Writer {
|
||||
return &logWriter{
|
||||
logger: l,
|
||||
}
|
||||
}
|
||||
|
||||
// RedirectStdLog redirects output from the standard library's package-global logger
|
||||
// to this logger at the specified level and with zero or more Field's. Since this logger already
|
||||
// handles caller annotations, timestamps, etc., it automatically disables the standard
|
||||
// library's annotations and prefixing.
|
||||
// A function is returned that restores the original prefix and flags and resets the standard
|
||||
// library's output to os.Stdout.
|
||||
func (l *Logger) RedirectStdLog(level Level, fields ...Field) func() {
|
||||
return l.log.Logr().RedirectStdLog(level, fields...)
|
||||
}
|
||||
|
||||
// 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 (l *Logger) RemoveTargets(ctx context.Context, f func(ti TargetInfo) bool) error {
|
||||
return l.log.Logr().RemoveTargets(ctx, f)
|
||||
}
|
||||
|
||||
// SetMetricsCollector sets (or resets) the metrics collector to be used for gathering
|
||||
// metrics for all targets. Only targets added after this call will use the collector.
|
||||
//
|
||||
// To ensure all targets use a collector, use the `SetMetricsCollector` option when
|
||||
// creating the Logger instead, or configure/reconfigure the Logger after calling this method.
|
||||
func (l *Logger) SetMetricsCollector(collector MetricsCollector, updateFrequencyMillis int64) {
|
||||
l.log.Logr().SetMetricsCollector(collector, updateFrequencyMillis)
|
||||
}
|
||||
|
||||
// Sugar creates a new `Logger` with a less structured API. Any fields are preserved.
|
||||
func (l *Logger) Sugar(fields ...Field) Sugar {
|
||||
return l.log.Sugar(fields...)
|
||||
}
|
||||
|
||||
// Flush forces all targets to write out any queued log records with a default timeout.
|
||||
func (l *Logger) Flush() error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), FlushTimeout)
|
||||
defer cancel()
|
||||
return l.log.Logr().FlushWithTimeout(ctx)
|
||||
}
|
||||
|
||||
// Flush forces all targets to write out any queued log records with the specfified timeout.
|
||||
func (l *Logger) FlushWithTimeout(ctx context.Context) error {
|
||||
return l.log.Logr().FlushWithTimeout(ctx)
|
||||
}
|
||||
|
||||
// Shutdown shuts down the logger after making best efforts to flush any
|
||||
// remaining records.
|
||||
func (l *Logger) Shutdown() error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), ShutdownTimeout)
|
||||
defer cancel()
|
||||
return l.log.Logr().ShutdownWithTimeout(ctx)
|
||||
}
|
||||
|
||||
// Shutdown shuts down the logger after making best efforts to flush any
|
||||
// remaining records.
|
||||
func (l *Logger) ShutdownWithTimeout(ctx context.Context) error {
|
||||
return l.log.Logr().ShutdownWithTimeout(ctx)
|
||||
}
|
||||
|
||||
// GetPackageName reduces a fully qualified function name to the package name
|
||||
// By sirupsen: https://github.com/sirupsen/logrus/blob/master/entry.go
|
||||
func GetPackageName(f string) string {
|
||||
for {
|
||||
lastPeriod := strings.LastIndex(f, ".")
|
||||
lastSlash := strings.LastIndex(f, "/")
|
||||
if lastPeriod > lastSlash {
|
||||
f = f[:lastPeriod]
|
||||
} else {
|
||||
break
|
||||
}
|
||||
}
|
||||
return f
|
||||
}
|
||||
|
||||
type logWriter struct {
|
||||
logger *Logger
|
||||
}
|
||||
|
||||
func (lw *logWriter) Write(p []byte) (int, error) {
|
||||
lw.logger.Info(string(p))
|
||||
return len(p), nil
|
||||
}
|
||||
|
||||
// ErrConfigurationLock is returned when one of a logger's configuration APIs is called
|
||||
// while the configuration is locked.
|
||||
var ErrConfigurationLock = errors.New("configuration is locked")
|
||||
55
shared/mlog/options.go
Обычный файл
55
shared/mlog/options.go
Обычный файл
@@ -0,0 +1,55 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import "github.com/mattermost/logr/v2"
|
||||
|
||||
// MaxQueueSize is the maximum number of log records that can be queued.
|
||||
// If exceeded, `OnQueueFull` is called which determines if the log
|
||||
// record will be dropped or block until add is successful.
|
||||
// Defaults to DefaultMaxQueueSize.
|
||||
func MaxQueueSize(size int) Option {
|
||||
return logr.MaxQueueSize(size)
|
||||
}
|
||||
|
||||
// OnLoggerError, when not nil, is called any time an internal
|
||||
// logging error occurs. For example, this can happen when a
|
||||
// target cannot connect to its data sink.
|
||||
func OnLoggerError(f func(error)) Option {
|
||||
return logr.OnLoggerError(f)
|
||||
}
|
||||
|
||||
// OnQueueFull, when not nil, is called on an attempt to add
|
||||
// a log record to a full Logr queue.
|
||||
// `MaxQueueSize` can be used to modify the maximum queue size.
|
||||
// This function should return quickly, with a bool indicating whether
|
||||
// the log record should be dropped (true) or block until the log record
|
||||
// is successfully added (false). If nil then blocking (false) is assumed.
|
||||
func OnQueueFull(f func(rec *LogRec, maxQueueSize int) bool) Option {
|
||||
return logr.OnQueueFull(f)
|
||||
}
|
||||
|
||||
// OnTargetQueueFull, when not nil, is called on an attempt to add
|
||||
// a log record to a full target queue provided the target supports reporting
|
||||
// this condition.
|
||||
// This function should return quickly, with a bool indicating whether
|
||||
// the log record should be dropped (true) or block until the log record
|
||||
// is successfully added (false). If nil then blocking (false) is assumed.
|
||||
func OnTargetQueueFull(f func(target Target, rec *LogRec, maxQueueSize int) bool) Option {
|
||||
return logr.OnTargetQueueFull(f)
|
||||
}
|
||||
|
||||
// SetMetricsCollector enables metrics collection by supplying a MetricsCollector.
|
||||
// The MetricsCollector provides counters and gauges that are updated by log targets.
|
||||
// `updateFreqMillis` determines how often polled metrics are updated. Defaults to 15000 (15 seconds)
|
||||
// and must be at least 250 so we don't peg the CPU.
|
||||
func SetMetricsCollector(collector MetricsCollector, updateFreqMillis int64) Option {
|
||||
return logr.SetMetricsCollector(collector, updateFreqMillis)
|
||||
}
|
||||
|
||||
// StackFilter provides a list of package names to exclude from the top of
|
||||
// stack traces. The Logr packages are automatically filtered.
|
||||
func StackFilter(pkg ...string) Option {
|
||||
return logr.StackFilter(pkg...)
|
||||
}
|
||||
@@ -1,87 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"strings"
|
||||
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
)
|
||||
|
||||
// Implementation of zapcore.Core to interpret log messages from a standard logger
|
||||
// and translate the levels to zapcore levels.
|
||||
type stdLogLevelInterpreterCore struct {
|
||||
wrappedCore zapcore.Core
|
||||
}
|
||||
|
||||
func stdLogInterpretZapEntry(entry zapcore.Entry) zapcore.Entry {
|
||||
message := entry.Message
|
||||
if strings.Index(message, "[DEBUG]") == 0 {
|
||||
entry.Level = zapcore.DebugLevel
|
||||
entry.Message = message[7:]
|
||||
} else if strings.Index(message, "[DEBG]") == 0 {
|
||||
entry.Level = zapcore.DebugLevel
|
||||
entry.Message = message[6:]
|
||||
} else if strings.Index(message, "[WARN]") == 0 {
|
||||
entry.Level = zapcore.WarnLevel
|
||||
entry.Message = message[6:]
|
||||
} else if strings.Index(message, "[ERROR]") == 0 {
|
||||
entry.Level = zapcore.ErrorLevel
|
||||
entry.Message = message[7:]
|
||||
} else if strings.Index(message, "[EROR]") == 0 {
|
||||
entry.Level = zapcore.ErrorLevel
|
||||
entry.Message = message[6:]
|
||||
} else if strings.Index(message, "[ERR]") == 0 {
|
||||
entry.Level = zapcore.ErrorLevel
|
||||
entry.Message = message[5:]
|
||||
} else if strings.Index(message, "[INFO]") == 0 {
|
||||
entry.Level = zapcore.InfoLevel
|
||||
entry.Message = message[6:]
|
||||
}
|
||||
return entry
|
||||
}
|
||||
|
||||
func (s *stdLogLevelInterpreterCore) Enabled(lvl zapcore.Level) bool {
|
||||
return s.wrappedCore.Enabled(lvl)
|
||||
}
|
||||
|
||||
func (s *stdLogLevelInterpreterCore) With(fields []zapcore.Field) zapcore.Core {
|
||||
return s.wrappedCore.With(fields)
|
||||
}
|
||||
|
||||
func (s *stdLogLevelInterpreterCore) Check(entry zapcore.Entry, checkedEntry *zapcore.CheckedEntry) *zapcore.CheckedEntry {
|
||||
entry = stdLogInterpretZapEntry(entry)
|
||||
return s.wrappedCore.Check(entry, checkedEntry)
|
||||
}
|
||||
|
||||
func (s *stdLogLevelInterpreterCore) Write(entry zapcore.Entry, fields []zapcore.Field) error {
|
||||
entry = stdLogInterpretZapEntry(entry)
|
||||
return s.wrappedCore.Write(entry, fields)
|
||||
}
|
||||
|
||||
func (s *stdLogLevelInterpreterCore) Sync() error {
|
||||
return s.wrappedCore.Sync()
|
||||
}
|
||||
|
||||
func getStdLogOption() zap.Option {
|
||||
return zap.WrapCore(
|
||||
func(core zapcore.Core) zapcore.Core {
|
||||
return &stdLogLevelInterpreterCore{core}
|
||||
},
|
||||
)
|
||||
}
|
||||
|
||||
type loggerWriter struct {
|
||||
logFunc func(msg string, fields ...Field)
|
||||
}
|
||||
|
||||
func (l *loggerWriter) Write(p []byte) (int, error) {
|
||||
trimmed := string(bytes.TrimSpace(p))
|
||||
for _, line := range strings.Split(trimmed, "\n") {
|
||||
l.logFunc(line)
|
||||
}
|
||||
return len(p), nil
|
||||
}
|
||||
@@ -1,42 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
"go.uber.org/zap/zapcore"
|
||||
)
|
||||
|
||||
func TestStdLogInterpretZapEntry(t *testing.T) {
|
||||
for _, tc := range []struct {
|
||||
testname string
|
||||
message string
|
||||
expectedMessage string
|
||||
expectedLevel zapcore.Level
|
||||
}{
|
||||
{"Debug Basic", "[DEBUG]My message", "My message", zapcore.DebugLevel},
|
||||
{"Debug Basic2", "[DEBG]My message", "My message", zapcore.DebugLevel},
|
||||
{"Warn Basic", "[WARN]My message", "My message", zapcore.WarnLevel},
|
||||
{"Error Basic", "[ERROR]My message", "My message", zapcore.ErrorLevel},
|
||||
{"Error Basic2", "[EROR]My message", "My message", zapcore.ErrorLevel},
|
||||
{"Error Basic3", "[ERR]My message", "My message", zapcore.ErrorLevel},
|
||||
{"Info Basic", "[INFO]My message", "My message", zapcore.InfoLevel},
|
||||
{"Unknown level", "[UNKNOWN]My message", "[UNKNOWN]My message", zapcore.PanicLevel},
|
||||
{"No level", "My message", "My message", zapcore.PanicLevel},
|
||||
{"Empty message", "", "", zapcore.PanicLevel},
|
||||
{"Malformed level", "INFO]My message", "INFO]My message", zapcore.PanicLevel},
|
||||
} {
|
||||
t.Run(tc.testname, func(t *testing.T) {
|
||||
inEntry := zapcore.Entry{
|
||||
Level: zapcore.PanicLevel,
|
||||
Message: tc.message,
|
||||
}
|
||||
resultEntry := stdLogInterpretZapEntry(inEntry)
|
||||
assert.Equal(t, tc.expectedMessage, resultEntry.Message)
|
||||
assert.Equal(t, tc.expectedLevel, resultEntry.Level)
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -1,30 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"go.uber.org/zap"
|
||||
)
|
||||
|
||||
// Made for the plugin interface, use the regular logger for other uses
|
||||
type SugarLogger struct {
|
||||
wrappedLogger *Logger
|
||||
zapSugar *zap.SugaredLogger
|
||||
}
|
||||
|
||||
func (l *SugarLogger) Debug(msg string, keyValuePairs ...interface{}) {
|
||||
l.zapSugar.Debugw(msg, keyValuePairs...)
|
||||
}
|
||||
|
||||
func (l *SugarLogger) Info(msg string, keyValuePairs ...interface{}) {
|
||||
l.zapSugar.Infow(msg, keyValuePairs...)
|
||||
}
|
||||
|
||||
func (l *SugarLogger) Error(msg string, keyValuePairs ...interface{}) {
|
||||
l.zapSugar.Errorw(msg, keyValuePairs...)
|
||||
}
|
||||
|
||||
func (l *SugarLogger) Warn(msg string, keyValuePairs ...interface{}) {
|
||||
l.zapSugar.Warnw(msg, keyValuePairs...)
|
||||
}
|
||||
@@ -1,142 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"crypto/x509"
|
||||
"encoding/base64"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
|
||||
"github.com/mattermost/logr"
|
||||
"github.com/wiggin77/merror"
|
||||
syslog "github.com/wiggin77/srslog"
|
||||
)
|
||||
|
||||
// Syslog outputs log records to local or remote syslog.
|
||||
type Syslog struct {
|
||||
logr.Basic
|
||||
w *syslog.Writer
|
||||
}
|
||||
|
||||
// SyslogParams provides parameters for dialing a syslog daemon.
|
||||
type SyslogParams struct {
|
||||
IP string `json:"IP"`
|
||||
Port int `json:"Port"`
|
||||
Tag string `json:"Tag"`
|
||||
TLS bool `json:"TLS"`
|
||||
Cert string `json:"Cert"`
|
||||
Insecure bool `json:"Insecure"`
|
||||
}
|
||||
|
||||
// NewSyslogTarget creates a target capable of outputting log records to remote or local syslog, with or without TLS.
|
||||
func NewSyslogTarget(filter logr.Filter, formatter logr.Formatter, params *SyslogParams, maxQueue int) (*Syslog, error) {
|
||||
network := "tcp"
|
||||
var config *tls.Config
|
||||
|
||||
if params.TLS {
|
||||
network = "tcp+tls"
|
||||
config = &tls.Config{InsecureSkipVerify: params.Insecure}
|
||||
if params.Cert != "" {
|
||||
pool, err := getCertPool(params.Cert)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
config.RootCAs = pool
|
||||
}
|
||||
}
|
||||
raddr := fmt.Sprintf("%s:%d", params.IP, params.Port)
|
||||
|
||||
writer, err := syslog.DialWithTLSConfig(network, raddr, syslog.LOG_INFO, params.Tag, config)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
s := &Syslog{w: writer}
|
||||
s.Basic.Start(s, s, filter, formatter, maxQueue)
|
||||
|
||||
return s, nil
|
||||
}
|
||||
|
||||
// Shutdown stops processing log records after making best effort to flush queue.
|
||||
func (s *Syslog) Shutdown(ctx context.Context) error {
|
||||
errs := merror.New()
|
||||
|
||||
err := s.Basic.Shutdown(ctx)
|
||||
errs.Append(err)
|
||||
|
||||
err = s.w.Close()
|
||||
errs.Append(err)
|
||||
|
||||
return errs.ErrorOrNil()
|
||||
}
|
||||
|
||||
// getCertPool returns a x509.CertPool containing the cert(s)
|
||||
// from `cert`, which can be a path to a .pem or .crt file,
|
||||
// or a base64 encoded cert.
|
||||
func getCertPool(cert string) (*x509.CertPool, error) {
|
||||
if cert == "" {
|
||||
return nil, errors.New("no cert provided")
|
||||
}
|
||||
|
||||
// first treat as a file and try to read.
|
||||
serverCert, err := ioutil.ReadFile(cert)
|
||||
if err != nil {
|
||||
// maybe it's a base64 encoded cert
|
||||
serverCert, err = base64.StdEncoding.DecodeString(cert)
|
||||
if err != nil {
|
||||
return nil, errors.New("cert cannot be read")
|
||||
}
|
||||
}
|
||||
|
||||
pool := x509.NewCertPool()
|
||||
if ok := pool.AppendCertsFromPEM(serverCert); ok {
|
||||
return pool, nil
|
||||
}
|
||||
return nil, errors.New("cannot parse cert")
|
||||
}
|
||||
|
||||
// Write converts the log record to bytes, via the Formatter,
|
||||
// and outputs to syslog.
|
||||
func (s *Syslog) Write(rec *logr.LogRec) error {
|
||||
_, stacktrace := s.IsLevelEnabled(rec.Level())
|
||||
|
||||
buf := rec.Logger().Logr().BorrowBuffer()
|
||||
defer rec.Logger().Logr().ReleaseBuffer(buf)
|
||||
|
||||
buf, err := s.Formatter().Format(rec, stacktrace, buf)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
txt := buf.String()
|
||||
|
||||
switch rec.Level() {
|
||||
case logr.Panic, logr.Fatal:
|
||||
err = s.w.Crit(txt)
|
||||
case logr.Error:
|
||||
err = s.w.Err(txt)
|
||||
case logr.Warn:
|
||||
err = s.w.Warning(txt)
|
||||
case logr.Debug, logr.Trace:
|
||||
err = s.w.Debug(txt)
|
||||
default:
|
||||
// logr.Info plus all custom levels.
|
||||
err = s.w.Info(txt)
|
||||
}
|
||||
|
||||
if err != nil {
|
||||
reporter := rec.Logger().Logr().ReportError
|
||||
reporter(fmt.Errorf("syslog write fail: %w", err))
|
||||
// syslog writer will try to reconnect.
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
// String returns a string representation of this target.
|
||||
func (s *Syslog) String() string {
|
||||
return "SyslogTarget"
|
||||
}
|
||||
@@ -1,38 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func Test_getCertPool(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
cert string
|
||||
wantErr bool
|
||||
}{
|
||||
{name: "garbage in, garbage out", wantErr: true, cert: "THISISNOTACERT"},
|
||||
{name: "good cert base64", wantErr: false, cert: "LS0tLS1CRUdJTiBDRVJUSUZJQ0FURS0tLS0tCk1JSURqekNDQW5lZ0F3SUJBZ0lSQVBZZlJTd2R6S29wQkt4WXhLcXNsSlV3RFFZSktvWklodmNOQVFFTEJRQXcKSnpFbE1DTUdBMVVFQXd3Y1RXRjBkR1Z5Ylc5emRDd2dTVzVqTGlCSmJuUmxjbTVoYkNCRFFUQWVGdzB4T1RBegpNakl3TURFME1UVmFGdzB5TWpBek1EWXdNREUwTVRWYU1Ec3hPVEEzQmdOVkJBTVRNRTFoZEhSbGNtMXZjM1FzCklFbHVZeTRnU1c1MFpYSnVZV3dnU1c1MFpYSnRaV1JwWVhSbElFRjFkR2h2Y21sMGVUQ0NBU0l3RFFZSktvWkkKaHZjTkFRRUJCUUFEZ2dFUEFEQ0NBUW9DZ2dFQkFNamxpUmRtdm5OTDR1L0pyL00yZFB3UW1USlhFQlkvVnE5UQp2QVU1MlgzdFJNQ1B4Y2FGeit4NmZ0dXZkTzJOZG9oWEdBbXR4OVFVNUxaY3ZGZVREcG9WRUJvOUErNGp0THZECkRaWWFUTkxwSm1vU29KSGFEYmRXWCtPQU9xeURpV1M3NDFMdWlNS1dIaGV3OVFPaXNhdDJaSU5QeGptQWQ5d0UKeHRoVE1nenN2N01VcW5NZXI4VTVPR1EwUXk3d0FtTlJjKzJLM3FQd2t4ZTJSVXZjdGU1MERVRk5neEVnaW5zaAp2cmtPWFIzODN2VUNaZnU3MnF1OG9nZ2ppUXB5VGxsdTVqZTJBcDZKTGpZTGtFTWlNcXJZQUR1V29yL1pId2E2CldyRnFWRVR4V2ZBVjV1OUVoMHdaTS9LS1l3UlF1dzl5K05hbnM3N0ZtVWwxdFZXV05OOENBd0VBQWFPQm9UQ0IKbmpBTUJnTlZIUk1FQlRBREFRSC9NQjBHQTFVZERnUVdCQlFZNFVxc3d5cjJoTy9IZXRadDJSRHhKZFRJUGpCaQpCZ05WSFNNRVd6QlpnQlJGWlhWZzJaNXROSXNXZVdqQkxFeTJ5ektiTUtFcnBDa3dKekVsTUNNR0ExVUVBd3djClRXRjBkR1Z5Ylc5emRDd2dTVzVqTGlCSmJuUmxjbTVoYkNCRFFZSVVFaWZHVU9NK2JJRlpvMXRralpCNVlHQnIKMHhFd0N3WURWUjBQQkFRREFnRUdNQTBHQ1NxR1NJYjNEUUVCQ3dVQUE0SUJBUUFFZGV4TDMwUTB6QkhtUEFIOApMaGRLN2RielcxQ21JTGJ4UlpsS0F3Uk4raEtSWGlNVzNNSElraE51b1Y5QWV2NjAyUStqYTRsV3NSaS9rdE9MCm5pMUZXeDVnU1NjZ2RHOEpHajQ3ZE9tb1QzdlhLWDcrdW1pdjRyUUxQRGw5L0RLTXV2MjA0T1lKcTZWVCt1TlUKNkM2a0wxNTdqR0pFTzc2SDRmTVo4b1lzRDdTcTB6amlOS3R1Q1lpaTBuZ0gzajNnQjFqQUNMcVJndmVVN01kVApwcU9WMktmWTMxK2g4VkJ0a1V2bGpOenRROXhOWThGam10MFNNZjdFM0ZhVWNhYXIzWkNyNzBHNWFVM2RLYmU3CjQ3dkdPQmE1dENxdzRZSzBqZ0RLaWQzSUpRdWw5YTNKMW1Tc0g4V3kzdG85Y0FWNEtHWkJRTG56Q1gxNWEvK3YKM3lWaAotLS0tLUVORCBDRVJUSUZJQ0FURS0tLS0tIAotLS0tLUJFR0lOIENFUlRJRklDQVRFLS0tLS0KTUlJRGZqQ0NBbWFnQXdJQkFnSVVFaWZHVU9NK2JJRlpvMXRralpCNVlHQnIweEV3RFFZSktvWklodmNOQVFFTApCUUF3SnpFbE1DTUdBMVVFQXd3Y1RXRjBkR1Z5Ylc5emRDd2dTVzVqTGlCSmJuUmxjbTVoYkNCRFFUQWVGdzB4Ck9UQXpNakV5TVRJNE5ETmFGdzB5T1RBek1UZ3lNVEk0TkROYU1DY3hKVEFqQmdOVkJBTU1IRTFoZEhSbGNtMXYKYzNRc0lFbHVZeTRnU1c1MFpYSnVZV3dnUTBFd2dnRWlNQTBHQ1NxR1NJYjNEUUVCQVFVQUE0SUJEd0F3Z2dFSwpBb0lCQVFESDBYcTVyTUJHcEtPVldUcGI1TW5hSklXRlAvdk90dkVrKzdoVnJmT2ZlMS81eDBLazNVZ0FIajg1Cm90YUVaRDFMaG4vSkxrRXFDaUUvVVhNSkZ3SkRsTmNPNENrZEtCU3BZWDRiS0FxeTVxL1gzUXdpb01TTnBKRzEKK1lZck5HQkgwc2dLY0tqeUNhTGhtcVlMRDB4WkRWT21XSVlCVTlqVVB5WHc1VTB0bnNWclRxR014VmttMXhDWQprckNXTjFab1VyTHZMME1DWmM1cXB4b1BUb3ByOVVPOWNxU0JTdXk2QlZXVnVFV0JaaHBxSHQrdWw4VnhoenpZCnExazRsN3IycXcrL3dtMWlKQmVkVGVCVmVXTmFnOEphVmZMZ3UrL1c3b0pWbFBPMzJQbzdwbnZIcDhpSjNiNEsKelh5VkhhVFg0UzZFbSs2TFY4ODU1VFlyU2h6bEFnTUJBQUdqZ2FFd2daNHdIUVlEVlIwT0JCWUVGRVZsZFdEWgpubTAwaXhaNWFNRXNUTGJMTXBzd01HSUdBMVVkSXdSYk1GbUFGRVZsZFdEWm5tMDBpeFo1YU1Fc1RMYkxNcHN3Cm9TdWtLVEFuTVNVd0l3WURWUVFEREJ4TllYUjBaWEp0YjNOMExDQkpibU11SUVsdWRHVnlibUZzSUVOQmdoUVMKSjhaUTR6NXNnVm1qVzJTTmtIbGdZR3ZURVRBTUJnTlZIUk1FQlRBREFRSC9NQXNHQTFVZER3UUVBd0lCQmpBTgpCZ2txaGtpRzl3MEJBUXNGQUFPQ0FRRUFQaUNXRm1vcHlBa1kyVDNaeW80eWFSUGhYMStWT1RNS0p0WTZFVWhxCi9HSHo2a3pFeXZDVUJmME44OTJjaWJHeGVrckVvSXRZOU5xTzZSUVJmb3dnK0duNWtjMTN6NE55TDJXOC9lb1QKWHkwWnZmYVFiVSsrZlE2cFZ0V3RNYmxETVU5eGlZZDcvTUR2SnBPMzI4bDFWaGNkcDhrRWkrbEN2cHkwc0NSYwpQeHpQaGJnQ01BYlpFR3grNFRNUWQ0U1pLemxSeFcvMmZmbHBSZWg2djFEdjBWRFVTWVFXd3NVbmFMcGRLSGZoCmE1azB2dXlTWWNzekU0WUtsWTB6YWtlRmxKZnA3ZkJwMXhUd2NkVzhhVGZ3MTVFaWNQTXdUYzZ4eEE0SkpVSngKY2RkdTgxN24xbmF5SzV1NnI5UWgxb0lWa3IwbkM5WUVMTU15NGRwUGdKODhTQT09Ci0tLS0tRU5EIENFUlRJRklDQVRFLS0tLS0K"},
|
||||
{name: "good cert file", wantErr: false, cert: "test-tls-client-cert.pem"},
|
||||
}
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
pool, err := getCertPool(tt.cert)
|
||||
if tt.wantErr {
|
||||
assert.Error(t, err)
|
||||
assert.Nil(t, pool)
|
||||
} else {
|
||||
assert.NoError(t, err)
|
||||
assert.NotNil(t, pool)
|
||||
|
||||
// Test PEM has 2 certs.
|
||||
subjects := pool.Subjects()
|
||||
assert.Len(t, subjects, 2)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
@@ -1,273 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/tls"
|
||||
"errors"
|
||||
"fmt"
|
||||
"net"
|
||||
_ "net/http/pprof"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/hashicorp/go-multierror"
|
||||
"github.com/mattermost/logr"
|
||||
)
|
||||
|
||||
const (
|
||||
DialTimeoutSecs = 30
|
||||
WriteTimeoutSecs = 30
|
||||
RetryBackoffMillis int64 = 100
|
||||
MaxRetryBackoffMillis int64 = 30 * 1000 // 30 seconds
|
||||
)
|
||||
|
||||
// TCP outputs log records to raw socket server.
|
||||
type TCP struct {
|
||||
logr.Basic
|
||||
|
||||
params *TCPParams
|
||||
addy string
|
||||
|
||||
mutex sync.Mutex
|
||||
conn net.Conn
|
||||
monitor chan struct{}
|
||||
shutdown chan struct{}
|
||||
}
|
||||
|
||||
// TCPParams provides parameters for dialing a socket server.
|
||||
type TCPParams struct {
|
||||
IP string `json:"IP"`
|
||||
Port int `json:"Port"`
|
||||
TLS bool `json:"TLS"`
|
||||
Cert string `json:"Cert"`
|
||||
Insecure bool `json:"Insecure"`
|
||||
}
|
||||
|
||||
// NewTPCTarget creates a target capable of outputting log records to a raw socket, with or without TLS.
|
||||
func NewTCPTarget(filter logr.Filter, formatter logr.Formatter, params *TCPParams, maxQueue int) (*TCP, error) {
|
||||
tcp := &TCP{
|
||||
params: params,
|
||||
addy: fmt.Sprintf("%s:%d", params.IP, params.Port),
|
||||
monitor: make(chan struct{}),
|
||||
shutdown: make(chan struct{}),
|
||||
}
|
||||
tcp.Basic.Start(tcp, tcp, filter, formatter, maxQueue)
|
||||
|
||||
return tcp, nil
|
||||
}
|
||||
|
||||
// getConn provides a net.Conn. If a connection already exists, it is returned immediately,
|
||||
// otherwise this method blocks until a new connection is created, timeout or shutdown.
|
||||
func (tcp *TCP) getConn() (net.Conn, error) {
|
||||
tcp.mutex.Lock()
|
||||
defer tcp.mutex.Unlock()
|
||||
|
||||
Log(LvlTCPLogTarget, "getConn enter", String("addy", tcp.addy))
|
||||
defer Log(LvlTCPLogTarget, "getConn exit", String("addy", tcp.addy))
|
||||
|
||||
if tcp.conn != nil {
|
||||
Log(LvlTCPLogTarget, "reusing existing conn", String("addy", tcp.addy)) // use "With" once Zap is removed
|
||||
return tcp.conn, nil
|
||||
}
|
||||
|
||||
type result struct {
|
||||
conn net.Conn
|
||||
err error
|
||||
}
|
||||
|
||||
connChan := make(chan result)
|
||||
ctx, cancel := context.WithTimeout(context.Background(), time.Second*DialTimeoutSecs)
|
||||
defer cancel()
|
||||
|
||||
go func(ctx context.Context, ch chan result) {
|
||||
Log(LvlTCPLogTarget, "dailing", String("addy", tcp.addy))
|
||||
conn, err := tcp.dial(ctx)
|
||||
if err == nil {
|
||||
tcp.conn = conn
|
||||
tcp.monitor = make(chan struct{})
|
||||
go monitor(tcp.conn, tcp.monitor, Log)
|
||||
}
|
||||
ch <- result{conn: conn, err: err}
|
||||
}(ctx, connChan)
|
||||
|
||||
select {
|
||||
case <-tcp.shutdown:
|
||||
return nil, errors.New("shutdown")
|
||||
case res := <-connChan:
|
||||
return res.conn, res.err
|
||||
}
|
||||
}
|
||||
|
||||
// dial connects to a TCP socket, and optionally performs a TLS handshake.
|
||||
// A non-nil context must be provided which can cancel the dial.
|
||||
func (tcp *TCP) dial(ctx context.Context) (net.Conn, error) {
|
||||
var dialer net.Dialer
|
||||
dialer.Timeout = time.Second * DialTimeoutSecs
|
||||
conn, err := dialer.DialContext(ctx, "tcp", fmt.Sprintf("%s:%d", tcp.params.IP, tcp.params.Port))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if !tcp.params.TLS {
|
||||
return conn, nil
|
||||
}
|
||||
|
||||
Log(LvlTCPLogTarget, "TLS handshake", String("addy", tcp.addy))
|
||||
|
||||
tlsconfig := &tls.Config{
|
||||
ServerName: tcp.params.IP,
|
||||
InsecureSkipVerify: tcp.params.Insecure,
|
||||
}
|
||||
if tcp.params.Cert != "" {
|
||||
pool, err := getCertPool(tcp.params.Cert)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
tlsconfig.RootCAs = pool
|
||||
}
|
||||
|
||||
tlsConn := tls.Client(conn, tlsconfig)
|
||||
if err := tlsConn.Handshake(); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return tlsConn, nil
|
||||
}
|
||||
|
||||
func (tcp *TCP) close() error {
|
||||
tcp.mutex.Lock()
|
||||
defer tcp.mutex.Unlock()
|
||||
|
||||
var err error
|
||||
if tcp.conn != nil {
|
||||
Log(LvlTCPLogTarget, "closing connection", String("addy", tcp.addy))
|
||||
close(tcp.monitor)
|
||||
err = tcp.conn.Close()
|
||||
tcp.conn = nil
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
// Shutdown stops processing log records after making best effort to flush queue.
|
||||
func (tcp *TCP) Shutdown(ctx context.Context) error {
|
||||
errs := &multierror.Error{}
|
||||
|
||||
Log(LvlTCPLogTarget, "shutting down", String("addy", tcp.addy))
|
||||
|
||||
if err := tcp.Basic.Shutdown(ctx); err != nil {
|
||||
errs = multierror.Append(errs, err)
|
||||
}
|
||||
|
||||
if err := tcp.close(); err != nil {
|
||||
errs = multierror.Append(errs, err)
|
||||
}
|
||||
|
||||
close(tcp.shutdown)
|
||||
return errs.ErrorOrNil()
|
||||
}
|
||||
|
||||
// Write converts the log record to bytes, via the Formatter, and outputs to the socket.
|
||||
// Called by dedicated target goroutine and will block until success or shutdown.
|
||||
func (tcp *TCP) Write(rec *logr.LogRec) error {
|
||||
_, stacktrace := tcp.IsLevelEnabled(rec.Level())
|
||||
|
||||
buf := rec.Logger().Logr().BorrowBuffer()
|
||||
defer rec.Logger().Logr().ReleaseBuffer(buf)
|
||||
|
||||
buf, err := tcp.Formatter().Format(rec, stacktrace, buf)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
try := 1
|
||||
backoff := RetryBackoffMillis
|
||||
for {
|
||||
select {
|
||||
case <-tcp.shutdown:
|
||||
return err
|
||||
default:
|
||||
}
|
||||
|
||||
conn, err := tcp.getConn()
|
||||
if err != nil {
|
||||
Log(LvlTCPLogTarget, "failed getting connection", String("addy", tcp.addy), Err(err))
|
||||
reporter := rec.Logger().Logr().ReportError
|
||||
reporter(fmt.Errorf("log target %s connection error: %w", tcp.String(), err))
|
||||
backoff = tcp.sleep(backoff)
|
||||
continue
|
||||
}
|
||||
|
||||
conn.SetWriteDeadline(time.Now().Add(time.Second * WriteTimeoutSecs))
|
||||
_, err = buf.WriteTo(conn)
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
|
||||
Log(LvlTCPLogTarget, "write error", String("addy", tcp.addy), Err(err))
|
||||
reporter := rec.Logger().Logr().ReportError
|
||||
reporter(fmt.Errorf("log target %s write error: %w", tcp.String(), err))
|
||||
|
||||
_ = tcp.close()
|
||||
|
||||
backoff = tcp.sleep(backoff)
|
||||
try++
|
||||
Log(LvlTCPLogTarget, "retrying write", String("addy", tcp.addy), Int("try", try))
|
||||
}
|
||||
}
|
||||
|
||||
// monitor continuously tries to read from the connection to detect socket close.
|
||||
// This is needed because TCP target uses a write only socket and Linux systems
|
||||
// take a long time to detect a loss of connectivity on a socket when only writing;
|
||||
// the writes simply fail without an error returned.
|
||||
func monitor(conn net.Conn, done <-chan struct{}, logFunc LogFuncCustom) {
|
||||
addy := conn.RemoteAddr().String()
|
||||
defer logFunc(LvlTCPLogTarget, "monitor exiting", String("addy", addy))
|
||||
|
||||
buf := make([]byte, 1)
|
||||
for {
|
||||
logFunc(LvlTCPLogTarget, "monitor loop", String("addy", addy))
|
||||
|
||||
select {
|
||||
case <-done:
|
||||
return
|
||||
case <-time.After(1 * time.Second):
|
||||
}
|
||||
|
||||
err := conn.SetReadDeadline(time.Now().Add(time.Second * 30))
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
_, err = conn.Read(buf)
|
||||
|
||||
if errt, ok := err.(net.Error); ok && errt.Timeout() {
|
||||
// read timeout is expected, keep looping.
|
||||
continue
|
||||
}
|
||||
|
||||
// Any other error closes the connection, forcing a reconnect.
|
||||
logFunc(LvlTCPLogTarget, "monitor closing connection", Err(err))
|
||||
conn.Close()
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
// String returns a string representation of this target.
|
||||
func (tcp *TCP) String() string {
|
||||
return fmt.Sprintf("TcpTarget[%s:%d]", tcp.params.IP, tcp.params.Port)
|
||||
}
|
||||
|
||||
func (tcp *TCP) sleep(backoff int64) int64 {
|
||||
select {
|
||||
case <-tcp.shutdown:
|
||||
case <-time.After(time.Millisecond * time.Duration(backoff)):
|
||||
}
|
||||
|
||||
nextBackoff := backoff + (backoff >> 1)
|
||||
if nextBackoff > MaxRetryBackoffMillis {
|
||||
nextBackoff = MaxRetryBackoffMillis
|
||||
}
|
||||
return nextBackoff
|
||||
}
|
||||
@@ -1,198 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"errors"
|
||||
"fmt"
|
||||
"io"
|
||||
"net"
|
||||
"sync"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/stretchr/testify/require"
|
||||
"github.com/wiggin77/merror"
|
||||
)
|
||||
|
||||
const (
|
||||
testPort = 18066
|
||||
)
|
||||
|
||||
func TestNewTCPTarget(t *testing.T) {
|
||||
target := LogTarget{
|
||||
Type: "tcp",
|
||||
Format: "json",
|
||||
Levels: []LogLevel{LvlInfo},
|
||||
Options: []byte(`{"IP": "localhost", "Port": 18066}`),
|
||||
MaxQueueSize: 1000,
|
||||
}
|
||||
targets := map[string]*LogTarget{"tcp_test": &target}
|
||||
|
||||
t.Run("logging", func(t *testing.T) {
|
||||
buf := &buffer{}
|
||||
server, err := newSocketServer(testPort, buf)
|
||||
require.NoError(t, err)
|
||||
|
||||
data := []string{"I drink your milkshake!", "We don't need no badges!", "You can't fight in here! This is the war room!"}
|
||||
|
||||
logger := newLogr()
|
||||
err = logrAddTargets(logger, targets)
|
||||
require.NoError(t, err)
|
||||
|
||||
for _, s := range data {
|
||||
logger.Info(s)
|
||||
}
|
||||
err = logger.Logr().Flush()
|
||||
require.NoError(t, err)
|
||||
err = logger.Logr().Shutdown()
|
||||
require.NoError(t, err)
|
||||
|
||||
err = server.waitForAnyConnection()
|
||||
require.NoError(t, err)
|
||||
|
||||
err = server.stopServer(true)
|
||||
require.NoError(t, err)
|
||||
|
||||
sdata := buf.String()
|
||||
for _, s := range data {
|
||||
require.Contains(t, sdata, s)
|
||||
}
|
||||
})
|
||||
}
|
||||
|
||||
// socketServer is a simple socket server used for testing TCP log targets.
|
||||
// Note: There is more synchronization here than normally needed to avoid flaky tests.
|
||||
// For example, it's possible for a unit test to create a socketServer, attempt
|
||||
// writing to it, and stop the socket server before "go ss.listen()" gets scheduled.
|
||||
type socketServer struct {
|
||||
listener net.Listener
|
||||
anyConn chan struct{}
|
||||
buf *buffer
|
||||
conns map[string]*socketServerConn
|
||||
mux sync.Mutex
|
||||
}
|
||||
|
||||
type socketServerConn struct {
|
||||
raddy string
|
||||
conn net.Conn
|
||||
done chan struct{}
|
||||
}
|
||||
|
||||
func newSocketServer(port int, buf *buffer) (*socketServer, error) {
|
||||
ss := &socketServer{
|
||||
buf: buf,
|
||||
conns: make(map[string]*socketServerConn),
|
||||
anyConn: make(chan struct{}),
|
||||
}
|
||||
|
||||
addy := fmt.Sprintf(":%d", port)
|
||||
l, err := net.Listen("tcp4", addy)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
ss.listener = l
|
||||
|
||||
go ss.listen()
|
||||
return ss, nil
|
||||
}
|
||||
|
||||
func (ss *socketServer) listen() {
|
||||
for {
|
||||
conn, err := ss.listener.Accept()
|
||||
if err != nil {
|
||||
return
|
||||
}
|
||||
sconn := &socketServerConn{raddy: conn.RemoteAddr().String(), conn: conn, done: make(chan struct{})}
|
||||
ss.registerConnection(sconn)
|
||||
go ss.handleConnection(sconn)
|
||||
}
|
||||
}
|
||||
|
||||
func (ss *socketServer) waitForAnyConnection() error {
|
||||
var err error
|
||||
select {
|
||||
case <-ss.anyConn:
|
||||
case <-time.After(5 * time.Second):
|
||||
err = errors.New("wait for any connection timed out")
|
||||
}
|
||||
return err
|
||||
}
|
||||
|
||||
func (ss *socketServer) handleConnection(sconn *socketServerConn) {
|
||||
close(ss.anyConn)
|
||||
defer ss.unregisterConnection(sconn)
|
||||
buf := make([]byte, 1024)
|
||||
|
||||
for {
|
||||
n, err := sconn.conn.Read(buf)
|
||||
if n > 0 {
|
||||
ss.buf.Write(buf[:n])
|
||||
}
|
||||
if err == io.EOF {
|
||||
ss.signalDone(sconn)
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (ss *socketServer) registerConnection(sconn *socketServerConn) {
|
||||
ss.mux.Lock()
|
||||
defer ss.mux.Unlock()
|
||||
ss.conns[sconn.raddy] = sconn
|
||||
}
|
||||
|
||||
func (ss *socketServer) unregisterConnection(sconn *socketServerConn) {
|
||||
ss.mux.Lock()
|
||||
defer ss.mux.Unlock()
|
||||
delete(ss.conns, sconn.raddy)
|
||||
}
|
||||
|
||||
func (ss *socketServer) signalDone(sconn *socketServerConn) {
|
||||
ss.mux.Lock()
|
||||
defer ss.mux.Unlock()
|
||||
close(sconn.done)
|
||||
}
|
||||
|
||||
func (ss *socketServer) stopServer(wait bool) error {
|
||||
errs := merror.New()
|
||||
ss.listener.Close()
|
||||
|
||||
ss.mux.Lock()
|
||||
// defensive copy; no more connections can be accepted so copy will stay current.
|
||||
conns := make(map[string]*socketServerConn, len(ss.conns))
|
||||
for k, v := range ss.conns {
|
||||
conns[k] = v
|
||||
}
|
||||
ss.mux.Unlock()
|
||||
|
||||
for _, sconn := range conns {
|
||||
if wait {
|
||||
select {
|
||||
case <-sconn.done:
|
||||
case <-time.After(time.Second * 5):
|
||||
errs.Append(errors.New("timed out"))
|
||||
}
|
||||
}
|
||||
}
|
||||
return errs.ErrorOrNil()
|
||||
}
|
||||
|
||||
type buffer struct {
|
||||
buf bytes.Buffer
|
||||
mux sync.Mutex
|
||||
}
|
||||
|
||||
func (b *buffer) Write(p []byte) (n int, err error) {
|
||||
b.mux.Lock()
|
||||
defer b.mux.Unlock()
|
||||
return b.buf.Write(p)
|
||||
}
|
||||
|
||||
func (b *buffer) String() string {
|
||||
b.mux.Lock()
|
||||
defer b.mux.Unlock()
|
||||
return b.buf.String()
|
||||
}
|
||||
@@ -1,43 +0,0 @@
|
||||
-----BEGIN CERTIFICATE-----
|
||||
MIIDjzCCAnegAwIBAgIRAPYfRSwdzKopBKxYxKqslJUwDQYJKoZIhvcNAQELBQAw
|
||||
JzElMCMGA1UEAwwcTWF0dGVybW9zdCwgSW5jLiBJbnRlcm5hbCBDQTAeFw0xOTAz
|
||||
MjIwMDE0MTVaFw0yMjAzMDYwMDE0MTVaMDsxOTA3BgNVBAMTME1hdHRlcm1vc3Qs
|
||||
IEluYy4gSW50ZXJuYWwgSW50ZXJtZWRpYXRlIEF1dGhvcml0eTCCASIwDQYJKoZI
|
||||
hvcNAQEBBQADggEPADCCAQoCggEBAMjliRdmvnNL4u/Jr/M2dPwQmTJXEBY/Vq9Q
|
||||
vAU52X3tRMCPxcaFz+x6ftuvdO2NdohXGAmtx9QU5LZcvFeTDpoVEBo9A+4jtLvD
|
||||
DZYaTNLpJmoSoJHaDbdWX+OAOqyDiWS741LuiMKWHhew9QOisat2ZINPxjmAd9wE
|
||||
xthTMgzsv7MUqnMer8U5OGQ0Qy7wAmNRc+2K3qPwkxe2RUvcte50DUFNgxEginsh
|
||||
vrkOXR383vUCZfu72qu8oggjiQpyTllu5je2Ap6JLjYLkEMiMqrYADuWor/ZHwa6
|
||||
WrFqVETxWfAV5u9Eh0wZM/KKYwRQuw9y+Nans77FmUl1tVWWNN8CAwEAAaOBoTCB
|
||||
njAMBgNVHRMEBTADAQH/MB0GA1UdDgQWBBQY4Uqswyr2hO/HetZt2RDxJdTIPjBi
|
||||
BgNVHSMEWzBZgBRFZXVg2Z5tNIsWeWjBLEy2yzKbMKErpCkwJzElMCMGA1UEAwwc
|
||||
TWF0dGVybW9zdCwgSW5jLiBJbnRlcm5hbCBDQYIUEifGUOM+bIFZo1tkjZB5YGBr
|
||||
0xEwCwYDVR0PBAQDAgEGMA0GCSqGSIb3DQEBCwUAA4IBAQAEdexL30Q0zBHmPAH8
|
||||
LhdK7dbzW1CmILbxRZlKAwRN+hKRXiMW3MHIkhNuoV9Aev602Q+ja4lWsRi/ktOL
|
||||
ni1FWx5gSScgdG8JGj47dOmoT3vXKX7+umiv4rQLPDl9/DKMuv204OYJq6VT+uNU
|
||||
6C6kL157jGJEO76H4fMZ8oYsD7Sq0zjiNKtuCYii0ngH3j3gB1jACLqRgveU7MdT
|
||||
pqOV2KfY31+h8VBtkUvljNztQ9xNY8Fjmt0SMf7E3FaUcaar3ZCr70G5aU3dKbe7
|
||||
47vGOBa5tCqw4YK0jgDKid3IJQul9a3J1mSsH8Wy3to9cAV4KGZBQLnzCX15a/+v
|
||||
3yVh
|
||||
-----END CERTIFICATE-----
|
||||
-----BEGIN CERTIFICATE-----
|
||||
MIIDfjCCAmagAwIBAgIUEifGUOM+bIFZo1tkjZB5YGBr0xEwDQYJKoZIhvcNAQEL
|
||||
BQAwJzElMCMGA1UEAwwcTWF0dGVybW9zdCwgSW5jLiBJbnRlcm5hbCBDQTAeFw0x
|
||||
OTAzMjEyMTI4NDNaFw0yOTAzMTgyMTI4NDNaMCcxJTAjBgNVBAMMHE1hdHRlcm1v
|
||||
c3QsIEluYy4gSW50ZXJuYWwgQ0EwggEiMA0GCSqGSIb3DQEBAQUAA4IBDwAwggEK
|
||||
AoIBAQDH0Xq5rMBGpKOVWTpb5MnaJIWFP/vOtvEk+7hVrfOfe1/5x0Kk3UgAHj85
|
||||
otaEZD1Lhn/JLkEqCiE/UXMJFwJDlNcO4CkdKBSpYX4bKAqy5q/X3QwioMSNpJG1
|
||||
+YYrNGBH0sgKcKjyCaLhmqYLD0xZDVOmWIYBU9jUPyXw5U0tnsVrTqGMxVkm1xCY
|
||||
krCWN1ZoUrLvL0MCZc5qpxoPTopr9UO9cqSBSuy6BVWVuEWBZhpqHt+ul8VxhzzY
|
||||
q1k4l7r2qw+/wm1iJBedTeBVeWNag8JaVfLgu+/W7oJVlPO32Po7pnvHp8iJ3b4K
|
||||
zXyVHaTX4S6Em+6LV8855TYrShzlAgMBAAGjgaEwgZ4wHQYDVR0OBBYEFEVldWDZ
|
||||
nm00ixZ5aMEsTLbLMpswMGIGA1UdIwRbMFmAFEVldWDZnm00ixZ5aMEsTLbLMpsw
|
||||
oSukKTAnMSUwIwYDVQQDDBxNYXR0ZXJtb3N0LCBJbmMuIEludGVybmFsIENBghQS
|
||||
J8ZQ4z5sgVmjW2SNkHlgYGvTETAMBgNVHRMEBTADAQH/MAsGA1UdDwQEAwIBBjAN
|
||||
BgkqhkiG9w0BAQsFAAOCAQEAPiCWFmopyAkY2T3Zyo4yaRPhX1+VOTMKJtY6EUhq
|
||||
/GHz6kzEyvCUBf0N892cibGxekrEoItY9NqO6RQRfowg+Gn5kc13z4NyL2W8/eoT
|
||||
Xy0ZvfaQbU++fQ6pVtWtMblDMU9xiYd7/MDvJpO328l1Vhcdp8kEi+lCvpy0sCRc
|
||||
PxzPhbgCMAbZEGx+4TMQd4SZKzlRxW/2fflpReh6v1Dv0VDUSYQWwsUnaLpdKHfh
|
||||
a5k0vuySYcszE4YKlY0zakeFlJfp7fBp1xTwcdW8aTfw15EicPMwTc6xxA4JJUJx
|
||||
cddu817n1nayK5u6r9Qh1oIVkr0nC9YELMMy4dpPgJ88SA==
|
||||
-----END CERTIFICATE-----
|
||||
@@ -1,46 +0,0 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"io"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
)
|
||||
|
||||
// testingWriter is an io.Writer that writes through t.Log
|
||||
type testingWriter struct {
|
||||
tb testing.TB
|
||||
}
|
||||
|
||||
func (tw *testingWriter) Write(b []byte) (int, error) {
|
||||
tw.tb.Log(strings.TrimSpace(string(b)))
|
||||
return len(b), nil
|
||||
}
|
||||
|
||||
// NewTestingLogger creates a Logger that proxies logs through a testing interface.
|
||||
// This allows tests that spin up App instances to avoid spewing logs unless the test fails or -verbose is specified.
|
||||
func NewTestingLogger(tb testing.TB, writer io.Writer) *Logger {
|
||||
logWriter := &testingWriter{tb}
|
||||
multiWriter := io.MultiWriter(logWriter, writer)
|
||||
logWriterSync := zapcore.AddSync(multiWriter)
|
||||
|
||||
testingLogger := &Logger{
|
||||
consoleLevel: zap.NewAtomicLevelAt(getZapLevel("debug")),
|
||||
fileLevel: zap.NewAtomicLevelAt(getZapLevel("info")),
|
||||
logrLogger: newLogr(),
|
||||
mutex: &sync.RWMutex{},
|
||||
}
|
||||
|
||||
logWriterCore := zapcore.NewCore(makeEncoder(true, false), zapcore.Lock(logWriterSync), testingLogger.consoleLevel)
|
||||
|
||||
testingLogger.zap = zap.New(logWriterCore,
|
||||
zap.AddCaller(),
|
||||
)
|
||||
return testingLogger
|
||||
}
|
||||
141
shared/mlog/tlog.go
Обычный файл
141
shared/mlog/tlog.go
Обычный файл
@@ -0,0 +1,141 @@
|
||||
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
|
||||
// See LICENSE.txt for license information.
|
||||
|
||||
package mlog
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"io"
|
||||
"os"
|
||||
"strings"
|
||||
"sync"
|
||||
"testing"
|
||||
|
||||
"github.com/mattermost/logr/v2"
|
||||
"github.com/mattermost/logr/v2/formatters"
|
||||
"github.com/mattermost/logr/v2/targets"
|
||||
)
|
||||
|
||||
// CreateTestLogger creates a logger for unit tests, using the `TB.Log`
|
||||
func CreateTestLogger(tb testing.TB, writer io.Writer, levels ...Level) *Logger {
|
||||
logger, _ := NewLogger()
|
||||
|
||||
filter := logr.NewCustomFilter(levels...)
|
||||
formatter := &formatters.Plain{}
|
||||
|
||||
if tb != nil {
|
||||
testtarget := newTestingTarget(tb)
|
||||
if err := logger.log.Logr().AddTarget(testtarget, "_testTB", filter, formatter, 1000); err != nil {
|
||||
tb.Fail()
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
if writer != nil {
|
||||
target := targets.NewWriterTarget(writer)
|
||||
if err := logger.log.Logr().AddTarget(target, "_testWriter", filter, formatter, 1000); err != nil {
|
||||
tb.Fail()
|
||||
return nil
|
||||
}
|
||||
}
|
||||
return logger
|
||||
}
|
||||
|
||||
func AddWriterTarget(logger *Logger, w io.Writer, useJSON bool, levels ...Level) error {
|
||||
filter := logr.NewCustomFilter(levels...)
|
||||
|
||||
var formatter logr.Formatter
|
||||
if useJSON {
|
||||
formatter = &formatters.JSON{EnableCaller: true}
|
||||
} else {
|
||||
formatter = &formatters.Plain{EnableCaller: true}
|
||||
}
|
||||
|
||||
target := targets.NewWriterTarget(w)
|
||||
return logger.log.Logr().AddTarget(target, "_testWriter", filter, formatter, 1000)
|
||||
}
|
||||
|
||||
// CreateConsoleTestLogger creates a logger for unit tests. Log records are output to `os.Stdout`.
|
||||
// Logs can also be mirrored to the optional `io.Writer`.
|
||||
func CreateConsoleTestLogger(useJSON bool, level Level) *Logger {
|
||||
logger, _ := NewLogger()
|
||||
|
||||
filter := logr.StdFilter{
|
||||
Lvl: level,
|
||||
Stacktrace: LvlPanic,
|
||||
}
|
||||
|
||||
var formatter logr.Formatter
|
||||
if useJSON {
|
||||
formatter = &formatters.JSON{EnableCaller: true}
|
||||
} else {
|
||||
formatter = &formatters.Plain{EnableCaller: true}
|
||||
}
|
||||
|
||||
target := targets.NewWriterTarget(os.Stdout)
|
||||
if err := logger.log.Logr().AddTarget(target, "_testcon", filter, formatter, 1000); err != nil {
|
||||
panic(err)
|
||||
}
|
||||
return logger
|
||||
}
|
||||
|
||||
// testingTarget is a simple log target that writes to the testing log.
|
||||
type testingTarget struct {
|
||||
mux sync.Mutex
|
||||
tb testing.TB
|
||||
}
|
||||
|
||||
func newTestingTarget(tb testing.TB) *testingTarget {
|
||||
return &testingTarget{
|
||||
tb: tb,
|
||||
}
|
||||
}
|
||||
|
||||
// Init is called once to initialize the target.
|
||||
func (tt *testingTarget) Init() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Write outputs bytes to this file target.
|
||||
func (tt *testingTarget) Write(p []byte, rec *logr.LogRec) (int, error) {
|
||||
tt.mux.Lock()
|
||||
defer tt.mux.Unlock()
|
||||
|
||||
if tt.tb != nil {
|
||||
tt.tb.Helper()
|
||||
tt.tb.Log(strings.TrimSpace(string(p)))
|
||||
}
|
||||
return len(p), nil
|
||||
}
|
||||
|
||||
// Shutdown is called once to free/close any resources.
|
||||
// Target queue is already drained when this is called.
|
||||
func (tt *testingTarget) Shutdown() error {
|
||||
tt.mux.Lock()
|
||||
defer tt.mux.Unlock()
|
||||
|
||||
tt.tb = nil
|
||||
return nil
|
||||
}
|
||||
|
||||
// Buffer provides a thread-safe buffer useful for logging to memory in unit tests.
|
||||
type Buffer struct {
|
||||
buf bytes.Buffer
|
||||
mux sync.Mutex
|
||||
}
|
||||
|
||||
func (b *Buffer) Read(p []byte) (n int, err error) {
|
||||
b.mux.Lock()
|
||||
defer b.mux.Unlock()
|
||||
return b.buf.Read(p)
|
||||
}
|
||||
func (b *Buffer) Write(p []byte) (n int, err error) {
|
||||
b.mux.Lock()
|
||||
defer b.mux.Unlock()
|
||||
return b.buf.Write(p)
|
||||
}
|
||||
func (b *Buffer) String() string {
|
||||
b.mux.Lock()
|
||||
defer b.mux.Unlock()
|
||||
return b.buf.String()
|
||||
}
|
||||
Ссылка в новой задаче
Block a user