Adds Advanced Logging to server. Advanced Logging is an optional logging capability that allows customers to send log records to any number of destinations.

Supported destinations:
- file
- syslog (with out without TLS)
- raw TCP socket (with out without TLS)

Allows developers to specify discrete log levels as well as the standard trace, debug, info, ... panic. Existing code and logging API usage is unchanged.

Log records are emitted asynchronously to reduce latency to the caller. Supports hot-reloading of logger config, including adding removing targets.

Advanced Logging is configured within config.json via "LogSettings.AdvancedLoggingConfig" which can contain a filespec to another config file, a database DSN, or JSON.
Этот коммит содержится в:
Doug Lauder
2020-07-15 14:40:36 -04:00
коммит произвёл GitHub
родитель 4ba6c35813
Коммит 90ff87a77f
53 изменённых файлов: 1442 добавлений и 82 удалений

Просмотреть файл

@@ -4,9 +4,13 @@
package mlog
import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"github.com/mattermost/logr"
)
// defaultLog manually encodes the log to STDERR, providing a basic, default logging implementation
@@ -49,3 +53,31 @@ 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(target 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")
}

30
mlog/errors.go Обычный файл
Просмотреть файл

@@ -0,0 +1,30 @@
// 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,6 +4,9 @@
package mlog
import (
"context"
"github.com/mattermost/logr"
"go.uber.org/zap"
"go.uber.org/zap/zapcore"
)
@@ -11,6 +14,10 @@ import (
var globalLogger *Logger
func InitGlobalLogger(logger *Logger) {
// Clean up previous instance.
if globalLogger != nil && globalLogger.logrLogger != nil {
globalLogger.logrLogger.Logr().Shutdown()
}
glob := *logger
glob.zap = glob.zap.WithOptions(zap.AddCallerSkip(1))
globalLogger = &glob
@@ -19,6 +26,12 @@ func InitGlobalLogger(logger *Logger) {
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
}
func RedirectStdLog(logger *Logger) {
@@ -26,6 +39,12 @@ func RedirectStdLog(logger *Logger) {
}
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
// DON'T USE THIS Modify the level on the app logger
func GloballyDisableDebugLogForTest() {
@@ -42,3 +61,10 @@ 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
var ConfigAdvancedLogging ConfigFunc = defaultAdvancedConfig
var ShutdownAdvancedLogging ShutdownFunc = defaultAdvancedShutdown
var AddTarget AddTargetFunc = defaultAddTarget

32
mlog/levels.go Обычный файл
Просмотреть файл

@@ -0,0 +1,32 @@
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
// See LICENSE.txt for license information.
package mlog
// Standard levels
var (
LvlPanic = LogLevel{ID: 0, Name: "panic"}
LvlFatal = LogLevel{ID: 1, Name: "fatal"}
LvlError = LogLevel{ID: 2, Name: "error"}
LvlWarn = LogLevel{ID: 3, Name: "warn"}
LvlInfo = LogLevel{ID: 4, Name: "info"}
LvlDebug = LogLevel{ID: 5, Name: "debug"}
LvlTrace = LogLevel{ID: 6, Name: "trace"}
// used only by the logger
LvlLogError = LogLevel{ID: 11, Name: "logerror"}
)
// Register custom (discrete) levels here...
// ! ID's must not exceed 32,768 !
var (
// used by the audit system
LvlAuditDebug = LogLevel{ID: 100, Name: "AuditDebug"}
LvlAuditError = LogLevel{ID: 101, Name: "AuditError"}
// used by the TCP log target
LvlTcpLogTarget = LogLevel{ID: 105, Name: "TcpLogTarget"}
)
// Combinations for LogM (log multi)
var (
MLvlExample = []LogLevel{LvlAuditDebug, LvlDebug}
)

Просмотреть файл

@@ -4,10 +4,12 @@
package mlog
import (
"context"
"io"
"log"
"os"
"github.com/mattermost/logr"
"go.uber.org/zap"
"go.uber.org/zap/zapcore"
"gopkg.in/natefinch/lumberjack.v2"
@@ -52,6 +54,7 @@ type Logger struct {
zap *zap.Logger
consoleLevel zap.AtomicLevel
fileLevel zap.AtomicLevel
logrLogger *logr.Logger
}
func getZapLevel(level string) zapcore.Level {
@@ -107,7 +110,6 @@ func NewLogger(config *LoggerConfiguration) *Logger {
logger.zap = zap.New(combinedCore,
zap.AddCaller(),
)
return logger
}
@@ -123,6 +125,10 @@ func (l *Logger) SetConsoleLevel(level string) {
func (l *Logger) With(fields ...Field) *Logger {
newlogger := *l
newlogger.zap = newlogger.zap.With(fields...)
if newlogger.logrLogger != nil {
ll := newlogger.logrLogger.WithFields(zapToLogr(fields))
newlogger.logrLogger = &ll
}
return &newlogger
}
@@ -161,20 +167,98 @@ func (l *Logger) Sugar() *SugarLogger {
func (l *Logger) Debug(message string, fields ...Field) {
l.zap.Debug(message, fields...)
if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Debug) {
l.logrLogger.WithFields(zapToLogr(fields)).Debug(message)
}
}
func (l *Logger) Info(message string, fields ...Field) {
l.zap.Info(message, fields...)
if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Info) {
l.logrLogger.WithFields(zapToLogr(fields)).Info(message)
}
}
func (l *Logger) Warn(message string, fields ...Field) {
l.zap.Warn(message, fields...)
if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Warn) {
l.logrLogger.WithFields(zapToLogr(fields)).Warn(message)
}
}
func (l *Logger) Error(message string, fields ...Field) {
l.zap.Error(message, fields...)
if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Error) {
l.logrLogger.WithFields(zapToLogr(fields)).Error(message)
}
}
func (l *Logger) Critical(message string, fields ...Field) {
l.zap.Error(message, fields...)
if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Error) {
l.logrLogger.WithFields(zapToLogr(fields)).Error(message)
}
}
func (l *Logger) Log(level LogLevel, message string, fields ...Field) {
if l.logrLogger != nil && isLevelEnabled(l.logrLogger, logr.Level(level)) {
l.logrLogger.WithFields(zapToLogr(fields)).Log(logr.Level(level), message)
}
}
func (l *Logger) LogM(levels []LogLevel, message string, fields ...Field) {
if l.logrLogger != nil {
var logger *logr.Logger
for _, lvl := range levels {
if isLevelEnabled(l.logrLogger, logr.Level(lvl)) {
// don't create logger with fields unless at least one level is active.
if logger == nil {
l := l.logrLogger.WithFields(zapToLogr(fields))
logger = &l
}
logger.Log(logr.Level(lvl), message)
}
}
}
}
func (l *Logger) Flush(cxt context.Context) error {
if l.logrLogger != nil {
return l.logrLogger.Logr().Flush() // TODO: use context when Logr lib supports it.
}
return nil
}
// ShutdownAdvancedLogging stops the logger from accepting new log records and tries to
// flush queues within the context timeout. Once complete all targets are shutdown
// and any resources released.
func (l *Logger) ShutdownAdvancedLogging(cxt context.Context) error {
var err error
if l.logrLogger != nil {
err = l.logrLogger.Logr().Shutdown() // TODO: use context when Logr lib supports it.
l.logrLogger = nil
}
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 l.logrLogger != nil {
if err := l.ShutdownAdvancedLogging(context.Background()); err != nil {
Error("error shutting down previous logger", Err(err))
}
}
logr, err := newLogr(targets)
l.logrLogger = logr
return err
}
// AddTarget adds a logr.Target to the advanced logger. This is the preferred method
// to add custom targets or provide configuration that cannot be expressed via a
//config source.
func (l *Logger) AddTarget(target logr.Target) error {
return l.logrLogger.Logr().AddTarget(target)
}

211
mlog/logr.go Обычный файл
Просмотреть файл

@@ -0,0 +1,211 @@
// 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".
Format string // one of "json", "plain"
Levels []LogLevel
Options json.RawMessage
MaxQueueSize int
}
type LogTargetCfg map[string]*LogTarget
type LogrCleanup func() error
func newLogr(targets LogTargetCfg) (*logr.Logger, error) {
var errs error
lgr := logr.Logr{}
lgr.OnExit = func(int) {}
lgr.OnPanic = func(interface{}) {}
lgr.OnLoggerError = onLoggerError
lgr.OnQueueFull = onQueueFull
lgr.OnTargetQueueFull = onTargetQueueFull
for name, t := range targets {
target, err := newLogrTarget(name, t)
if err != nil {
errs = multierror.Append(err)
continue
}
lgr.AddTarget(target)
}
logger := lgr.NewLogger()
return &logger, errs
}
func newLogrTarget(name string, t *LogTarget) (logr.Target, error) {
formatter, err := newFormatter(name, t.Format)
if err != nil {
return nil, err
}
filter, err := newFilter(name, t.Levels)
if err != nil {
return nil, err
}
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)
}
return nil, fmt.Errorf("invalid type '%s' for target %s", t.Type, name)
}
func newFilter(name string, levels []LogLevel) (logr.Filter, error) {
filter := &logr.CustomFilter{}
for _, lvl := range levels {
filter.Add(logr.Level(lvl))
}
return filter, nil
}
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
}
if options.Filename == "" {
return nil, fmt.Errorf("missing 'Filename' option for target %s", name)
}
if err := checkFileWritable(options.Filename); err != nil {
return nil, fmt.Errorf("error writing to 'Filename' for target %s: %w", name, err)
}
newTarget := target.NewFileTarget(filter, formatter, target.FileOptions(*options), t.MaxQueueSize)
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 {
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)
}

142
mlog/syslog.go Обычный файл
Просмотреть файл

@@ -0,0 +1,142 @@
// 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"
}

38
mlog/syslog_test.go Обычный файл
Просмотреть файл

@@ -0,0 +1,38 @@
// 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)
}
})
}
}

274
mlog/tcp.go Обычный файл
Просмотреть файл

@@ -0,0 +1,274 @@
// Copyright (c) 2015-present Mattermost, Inc. All Rights Reserved.
// See LICENSE.txt for license information.
package mlog
import (
"context"
"crypto/tls"
"errors"
"fmt"
"net"
"sync"
"time"
"github.com/hashicorp/go-multierror"
"github.com/mattermost/logr"
_ "net/http/pprof"
)
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"`
}
// NewTcpTarget 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)
}
connChan <- 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{}) {
addy := conn.RemoteAddr().String()
defer Log(LvlTcpLogTarget, "monitor exiting", String("addy", addy))
buf := make([]byte, 1)
for {
Log(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.
Log(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
}

197
mlog/tcp_test.go Обычный файл
Просмотреть файл

@@ -0,0 +1,197 @@
// 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!"}
logr, err := newLogr(targets)
require.NoError(t, err)
for _, s := range data {
logr.Info(s)
}
err = logr.Logr().Flush()
require.NoError(t, err)
err = logr.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()
}

43
mlog/test-tls-client-cert.pem Обычный файл
Просмотреть файл

@@ -0,0 +1,43 @@
-----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-----