diff --git a/.changeset/logger-with-tee.md b/.changeset/logger-with-tee.md new file mode 100644 index 000000000..28e9d6558 --- /dev/null +++ b/.changeset/logger-with-tee.md @@ -0,0 +1,8 @@ +--- +"github.com/livekit/protocol": minor +"@livekit/protocol": patch +--- + +Replace logger.WithTap with logger.WithTee, which duplicates every log entry to a caller-supplied zaputil.Tee. The tee's core is built from the level each derived logger resolves, so the copy follows component levels rather than carrying a level of its own. + +Replace ZapLogger.WithMinLevel with WithComponentLeveler, which attaches a zaputil.ComponentLeveler to a branch of the logger tree. A leveler owns the per-component level and write-enabler cache for one configuration source and can only widen its parent, so a caller can resolve levels per (tenant, component) without rebuilding loggers on a config change. ZapLogger.Leveler exposes the leveler a branch resolves through, for use as the parent of a derived one. diff --git a/livekit/egress_test.go b/livekit/egress_test.go index 0d3d1b549..82dae8fa1 100644 --- a/livekit/egress_test.go +++ b/livekit/egress_test.go @@ -18,7 +18,6 @@ import ( "testing" "github.com/stretchr/testify/require" - "go.uber.org/zap/zapcore" "github.com/livekit/protocol/logger" "github.com/livekit/protocol/logger/testutil" @@ -32,7 +31,7 @@ type TestEgressLogOutput struct { func TestLoggerProto(t *testing.T) { ws := &testutil.BufferedWriteSyncer{} - l, err := logger.NewZapLogger(&logger.Config{}, logger.WithTap(zaputil.NewWriteEnabler(ws, zapcore.DebugLevel))) + l, err := logger.NewZapLogger(&logger.Config{Level: "debug"}, logger.WithTee(zaputil.NewTee(testutil.NewJSONCoreFactory(ws)))) require.NoError(t, err) s3 := &S3Upload{ diff --git a/logger/config.go b/logger/config.go index 1d1e7dee0..75c77477e 100644 --- a/logger/config.go +++ b/logger/config.go @@ -14,7 +14,12 @@ package logger -import "sync" +import ( + "strings" + "sync" + + "go.uber.org/zap/zapcore" +) type Config struct { JSON bool `yaml:"json,omitempty"` @@ -72,3 +77,18 @@ func (c *Config) AddUpdateObserver(cb ConfigObserver) { defer c.lock.Unlock() c.onUpdatedCallbacks = append(c.onUpdatedCallbacks, cb) } + +// ResolveComponentLevel always resolves: an unconfigured component takes Level. +func (c *Config) ResolveComponentLevel(component string) (zapcore.Level, bool) { + c.lock.Lock() + defer c.lock.Unlock() + + parts := strings.Split(component, ".") + for len(parts)> 0 { + if lvl, ok := c.ComponentLevels[strings.Join(parts, ".")]; ok { + return ParseZapLevel(lvl), true + } + parts = parts[:len(parts)-1] + } + return ParseZapLevel(c.Level), true +} diff --git a/logger/config_test.go b/logger/config_test.go new file mode 100644 index 000000000..55fb4054f --- /dev/null +++ b/logger/config_test.go @@ -0,0 +1,25 @@ +package logger + +import ( + "io" + "testing" + + "github.com/stretchr/testify/require" + "go.uber.org/zap/zapcore" + + "github.com/livekit/protocol/logger/zaputil" +) + +func TestConfigResolveComponentLevel(t *testing.T) { + conf := &Config{Level: "info", ComponentLevels: map[string]string{"rtc.room": "debug"}} + root := zaputil.NewRootComponentLeveler(zapcore.AddSync(io.Discard), conf) + + require.True(t, root.ComponentLevel("rtc.room").Enabled(zapcore.DebugLevel)) + require.True(t, root.ComponentLevel("rtc.room.track").Enabled(zapcore.DebugLevel)) + require.False(t, root.ComponentLevel("rtc").Enabled(zapcore.DebugLevel)) + require.True(t, root.ComponentLevel("rtc").Enabled(zapcore.InfoLevel)) + + lvl, ok := (&Config{}).ResolveComponentLevel("anything") + require.True(t, ok) + require.Equal(t, zapcore.InfoLevel, lvl) +} diff --git a/logger/logger.go b/logger/logger.go index c60d030a3..0d1af7981 100644 --- a/logger/logger.go +++ b/logger/logger.go @@ -18,14 +18,11 @@ import ( "log/slog" "os" "slices" - "strings" - "sync" "sync/atomic" "time" "github.com/go-logr/logr" "github.com/go-logr/logr/funcr" - "github.com/puzpuzpuz/xsync/v4" "go.uber.org/zap" "go.uber.org/zap/zapcore" @@ -88,7 +85,12 @@ func ParseZapLevel(level string) zapcore.Level { return lvl } -type DeferredFieldResolver = zaputil.DeferredFieldResolver +type ( + ComponentLeveler = zaputil.ComponentLeveler + ComponentLevelResolver = zaputil.ComponentLevelResolver + DeferredFieldResolver = zaputil.DeferredFieldResolver + FixedComponentLevel = zaputil.FixedComponentLevel +) type Logger interface { Debugw(msg string, keysAndValues ...any) @@ -140,86 +142,18 @@ func (l UnlikelyLogger) WithValues(keysAndValues ...any) UnlikelyLogger { return UnlikelyLogger{l.logger, slices.Concat(l.keysAndValues, keysAndValues)} } -type sharedConfig struct { - level zap.AtomicLevel - mu sync.Mutex - componentLevels map[string]zap.AtomicLevel - config *Config -} - -func newSharedConfig(conf *Config) *sharedConfig { - sc := &sharedConfig{ - level: zap.NewAtomicLevelAt(ParseZapLevel(conf.Level)), - config: conf, - componentLevels: make(map[string]zap.AtomicLevel), - } - conf.AddUpdateObserver(sc.onConfigUpdate) - _ = sc.onConfigUpdate(conf) - return sc -} - -func (c *sharedConfig) onConfigUpdate(conf *Config) error { - // update log levels - c.level.SetLevel(ParseZapLevel(conf.Level)) - - // we have to update alla existing component levels - c.mu.Lock() - c.config = conf - for component, atomicLevel := range c.componentLevels { - effectiveLevel := c.level.Level() - parts := strings.Split(component, ".") - confSearch: - for len(parts)> 0 { - search := strings.Join(parts, ".") - if compLevel, ok := conf.ComponentLevels[search]; ok { - effectiveLevel = ParseZapLevel(compLevel) - break confSearch - } - parts = parts[:len(parts)-1] - } - atomicLevel.SetLevel(effectiveLevel) - } - c.mu.Unlock() - return nil -} - -// ensure we have an atomic level in the map representing the full component path -// this makes it possible to update the log level after the fact -func (c *sharedConfig) ComponentLevel(component string) zap.AtomicLevel { - c.mu.Lock() - defer c.mu.Unlock() - if compLevel, ok := c.componentLevels[component]; ok { - return compLevel - } - - // search up the hierarchy to find the first level that is set - atomicLevel := zap.NewAtomicLevelAt(c.level.Level()) - c.componentLevels[component] = atomicLevel - parts := strings.Split(component, ".") - for len(parts)> 0 { - search := strings.Join(parts, ".") - if compLevel, ok := c.config.ComponentLevels[search]; ok { - atomicLevel.SetLevel(ParseZapLevel(compLevel)) - return atomicLevel - } - parts = parts[:len(parts)-1] - } - return atomicLevel -} - type zapConfig struct { - conf *Config - sc *sharedConfig - writeEnablers *xsync.Map[string, *zaputil.WriteEnabler] - levelEnablers *xsync.Map[string, *zaputil.OrLevelEnabler] - tap *zaputil.WriteEnabler + conf *Config + root *ComponentLeveler + tee zaputil.Tee } type ZapLoggerOption func(*zapConfig) -func WithTap(tap *zaputil.WriteEnabler) ZapLoggerOption { +// The tee does its own encoding, and shares the console's resolved level. +func WithTee(tee zaputil.Tee) ZapLoggerOption { return func(zc *zapConfig) { - zc.tap = tap + zc.tee = tee } } @@ -230,18 +164,24 @@ type ZapComponentLeveler interface { type ZapLogger interface { Logger ToZap() *zap.SugaredLogger + // ComponentLeveler names components relative to this logger, unlike Leveler, which + // takes paths from the root. ComponentLeveler() ZapComponentLeveler - WithMinLevel(lvl zapcore.LevelEnabler) Logger + Leveler() *ComponentLeveler + // lv supplies the write syncer as well as the level, so it must be derived from this + // logger's Leveler or output goes wherever its root points. + WithComponentLeveler(lv *ComponentLeveler) Logger } -type zapLogger[T zaputil.Encoder[T]] struct { +type zapLogger struct { zap *zap.SugaredLogger *zapConfig - enc T + enc zaputil.Encoder component string deferred []*zaputil.Deferrer sampler *zaputil.Sampler - minLevel zapcore.LevelEnabler + leveler *ComponentLeveler + tee zaputil.Tee } func FromZapLogger(log *zap.Logger, conf *Config, opts ...ZapLoggerOption) (ZapLogger, error) { @@ -251,12 +191,13 @@ func FromZapLogger(log *zap.Logger, conf *Config, opts ...ZapLoggerOption) (ZapL zap := log.WithOptions(zap.AddCallerSkip(1)).Sugar() zc := &zapConfig{ - conf: conf, - sc: newSharedConfig(conf), - writeEnablers: xsync.NewMap[string, *zaputil.WriteEnabler](), - levelEnablers: xsync.NewMap[string, *zaputil.OrLevelEnabler](), - tap: zaputil.NewDiscardWriteEnabler(), + conf: conf, + root: zaputil.NewRootComponentLeveler(os.Stderr, conf), } + conf.AddUpdateObserver(func(*Config) error { + zc.root.Refresh() + return nil + }) for _, opt := range opts { opt(zc) } @@ -276,38 +217,34 @@ func FromZapLogger(log *zap.Logger, conf *Config, opts ...ZapLoggerOption) (ZapL if conf.JSON { return newZapLogger(zap, zc, zaputil.NewProductionEncoder(), sampler), nil - } else { - return newZapLogger(zap, zc, zaputil.NewDevelopmentEncoder(), sampler), nil } + return newZapLogger(zap, zc, zaputil.NewDevelopmentEncoder(), sampler), nil } func NewZapLogger(conf *Config, opts ...ZapLoggerOption) (ZapLogger, error) { return FromZapLogger(nil, conf, opts...) } -func newZapLogger[T zaputil.Encoder[T]](zap *zap.SugaredLogger, zc *zapConfig, enc T, sampler *zaputil.Sampler) ZapLogger { - l := &zapLogger[T]{ +func newZapLogger(zap *zap.SugaredLogger, zc *zapConfig, enc zaputil.Encoder, sampler *zaputil.Sampler) ZapLogger { + l := &zapLogger{ zap: zap, zapConfig: zc, enc: enc, sampler: sampler, + leveler: zc.root, + tee: zc.tee, } l.zap = l.makeZap() return l } -func (l *zapLogger[T]) makeZap() *zap.SugaredLogger { - var console *zaputil.WriteEnabler - if l.minLevel == nil { - console, _ = l.writeEnablers.LoadOrCompute(l.component, func() (*zaputil.WriteEnabler, bool) { - return zaputil.NewWriteEnabler(os.Stderr, l.sc.ComponentLevel(l.component)), false - }) - } else { - enab := zaputil.OrLevelEnabler{l.minLevel, l.sc.ComponentLevel(l.component)} - console = zaputil.NewWriteEnabler(os.Stderr, enab) - } +func (l *zapLogger) makeZap() *zap.SugaredLogger { + console := l.leveler.WriteEnabler(l.component) - c := l.enc.Core(console, l.tap) + c := l.enc.Core(console) + if tee := l.tee.Core(console); tee != nil { + c = zapcore.NewTee(c, tee) + } for i := range l.deferred { c = zaputil.NewDeferredValueCore(c, l.deferred[i]) } @@ -318,76 +255,81 @@ func (l *zapLogger[T]) makeZap() *zap.SugaredLogger { return l.zap.WithOptions(zap.WrapCore(func(zapcore.Core) zapcore.Core { return c })) } -func (l *zapLogger[T]) ToZap() *zap.SugaredLogger { +func (l *zapLogger) ToZap() *zap.SugaredLogger { return l.zap.WithOptions(zap.AddCallerSkip(-1)) } -type zapLoggerComponentLeveler[T zaputil.Encoder[T]] struct { - zl *zapLogger[T] +type zapLoggerComponentLeveler struct { + zl *zapLogger } -func (l zapLoggerComponentLeveler[T]) ComponentLevel(component string) zapcore.LevelEnabler { +func (l zapLoggerComponentLeveler) ComponentLevel(component string) zapcore.LevelEnabler { if l.zl.component != "" { component = l.zl.component + "." + component } - enab, _ := l.zl.levelEnablers.LoadOrCompute(component, func() (*zaputil.OrLevelEnabler, bool) { - return &zaputil.OrLevelEnabler{l.zl.sc.ComponentLevel(component), l.zl.tap}, false - }) - return enab + return l.zl.leveler.ComponentLevel(component) } -func (l *zapLogger[T]) ComponentLeveler() ZapComponentLeveler { - return zapLoggerComponentLeveler[T]{l} +func (l *zapLogger) ComponentLeveler() ZapComponentLeveler { + return zapLoggerComponentLeveler{l} } -func (l *zapLogger[T]) Debugw(msg string, keysAndValues ...any) { +func (l *zapLogger) Debugw(msg string, keysAndValues ...any) { l.zap.Debugw(msg, keysAndValues...) } -func (l *zapLogger[T]) WithMinLevel(lvl zapcore.LevelEnabler) Logger { +func (l *zapLogger) Leveler() *ComponentLeveler { + return l.leveler +} + +func (l *zapLogger) WithComponentLeveler(lv *ComponentLeveler) Logger { + if lv == nil { + return l + } dup := *l - dup.minLevel = lvl + dup.leveler = lv dup.zap = dup.makeZap() return &dup } -func (l *zapLogger[T]) Infow(msg string, keysAndValues ...any) { +func (l *zapLogger) Infow(msg string, keysAndValues ...any) { l.zap.Infow(msg, keysAndValues...) } -func (l *zapLogger[T]) Warnw(msg string, err error, keysAndValues ...any) { +func (l *zapLogger) Warnw(msg string, err error, keysAndValues ...any) { if err != nil { keysAndValues = append(keysAndValues, "error", err) } l.zap.Warnw(msg, keysAndValues...) } -func (l *zapLogger[T]) Errorw(msg string, err error, keysAndValues ...any) { +func (l *zapLogger) Errorw(msg string, err error, keysAndValues ...any) { if err != nil { keysAndValues = append(keysAndValues, "error", err) } l.zap.Errorw(msg, keysAndValues...) } -func (l *zapLogger[T]) WithValues(keysAndValues ...any) Logger { +func (l *zapLogger) WithValues(keysAndValues ...any) Logger { dup := *l dup.enc = dup.enc.WithValues(keysAndValues...) + dup.tee = dup.tee.WithValues(keysAndValues...) dup.zap = dup.makeZap() return &dup } -func (l *zapLogger[T]) WithUnlikelyValues(keysAndValues ...any) UnlikelyLogger { +func (l *zapLogger) WithUnlikelyValues(keysAndValues ...any) UnlikelyLogger { return UnlikelyLogger{l, keysAndValues} } -func (l *zapLogger[T]) WithName(name string) Logger { +func (l *zapLogger) WithName(name string) Logger { dup := *l dup.zap = dup.zap.Named(name) return &dup } -func (l *zapLogger[T]) WithComponent(component string) Logger { +func (l *zapLogger) WithComponent(component string) Logger { dup := *l dup.zap = dup.zap.Named(component) if dup.component == "" { @@ -399,13 +341,13 @@ func (l *zapLogger[T]) WithComponent(component string) Logger { return &dup } -func (l *zapLogger[T]) WithCallDepth(depth int) Logger { +func (l *zapLogger) WithCallDepth(depth int) Logger { dup := *l dup.zap = dup.zap.WithOptions(zap.AddCallerSkip(depth)) return &dup } -func (l *zapLogger[T]) WithItemSampler() Logger { +func (l *zapLogger) WithItemSampler() Logger { if l.conf.ItemSampleSeconds == 0 { return l } @@ -419,14 +361,14 @@ func (l *zapLogger[T]) WithItemSampler() Logger { return &dup } -func (l *zapLogger[T]) WithoutSampler() Logger { +func (l *zapLogger) WithoutSampler() Logger { dup := *l dup.sampler = nil dup.zap = dup.makeZap() return &dup } -func (l *zapLogger[T]) WithDeferredValues() (Logger, DeferredFieldResolver) { +func (l *zapLogger) WithDeferredValues() (Logger, DeferredFieldResolver) { dup := *l def := &zaputil.Deferrer{} dup.deferred = append(dup.deferred[0:len(dup.deferred):len(dup.deferred)], def) diff --git a/logger/logger_bench_test.go b/logger/logger_bench_test.go new file mode 100644 index 000000000..4652c5304 --- /dev/null +++ b/logger/logger_bench_test.go @@ -0,0 +1,166 @@ +package logger + +import ( + "io" + "testing" + + "go.uber.org/zap/zapcore" + + "github.com/livekit/protocol/logger/testutil" + "github.com/livekit/protocol/logger/zaputil" + "github.com/livekit/protocol/utils/must" +) + +var benchSink Logger + +type loggerBenchCase struct { + label string + conf *Config + opts []ZapLoggerOption +} + +func (c loggerBenchCase) logger(b *testing.B) ZapLogger { + b.Helper() + return must.Get(NewZapLogger(c.conf, c.opts...)) +} + +func derivationBenchCases() []loggerBenchCase { + return []loggerBenchCase{ + {label: "console", conf: &Config{}}, + {label: "json", conf: &Config{JSON: true}}, + {label: "console tee", conf: &Config{}, opts: withDiscardTee()}, + {label: "json tee", conf: &Config{JSON: true}, opts: withDiscardTee()}, + } +} + +func withDiscardTee() []ZapLoggerOption { + return []ZapLoggerOption{WithTee(zaputil.NewTee(testutil.NewJSONCoreFactory(zapcore.AddSync(io.Discard))))} +} + +func BenchmarkLoggerWithComponent(b *testing.B) { + for _, c := range derivationBenchCases() { + b.Run(c.label, func(b *testing.B) { + root := c.logger(b) + branch := root.WithComponentLeveler(zaputil.NewComponentLeveler(root.Leveler(), FixedComponentLevel(zapcore.DebugLevel))) + + // One repeated component, so this measures the cache rather than resolution. + b.Run("root", func(b *testing.B) { + b.ReportAllocs() + for range b.N { + benchSink = root.WithComponent("rtc") + } + }) + + b.Run("branch leveler", func(b *testing.B) { + b.ReportAllocs() + for range b.N { + benchSink = branch.WithComponent("rtc") + } + }) + }) + } +} + +func BenchmarkLoggerWithComponentLeveler(b *testing.B) { + for _, c := range derivationBenchCases() { + b.Run(c.label, func(b *testing.B) { + root := c.logger(b) + lv := zaputil.NewComponentLeveler(root.Leveler(), FixedComponentLevel(zapcore.DebugLevel)) + + b.ReportAllocs() + for range b.N { + benchSink = root.WithComponentLeveler(lv) + } + }) + } +} + +func BenchmarkLoggerWithValues(b *testing.B) { + for _, c := range derivationBenchCases() { + b.Run(c.label, func(b *testing.B) { + root := c.logger(b) + + b.Run("root", func(b *testing.B) { + b.ReportAllocs() + for range b.N { + benchSink = root.WithValues("participant", "PA_1234") + } + }) + + b.Run("derived", func(b *testing.B) { + derived := root.WithValues("room", "RM_1234") + b.ReportAllocs() + for range b.N { + benchSink = derived.WithValues("participant", "PA_1234") + } + }) + }) + } +} + +func BenchmarkLoggerWithDeferredValues(b *testing.B) { + for _, c := range derivationBenchCases() { + b.Run(c.label, func(b *testing.B) { + root := c.logger(b) + + b.Run("derive", func(b *testing.B) { + b.ReportAllocs() + for range b.N { + benchSink, _ = root.WithDeferredValues() + } + }) + + b.Run("derive and resolve", func(b *testing.B) { + b.ReportAllocs() + for range b.N { + l, resolver := root.WithDeferredValues() + resolver.Resolve("participant", "PA_1234") + benchSink = l + } + }) + }) + } +} + +func BenchmarkLoggerWrite(b *testing.B) { + silenceStderr(b) + + cases := []loggerBenchCase{ + {label: "console", conf: &Config{Level: "debug"}}, + {label: "json", conf: &Config{JSON: true, Level: "debug"}}, + {label: "level disabled", conf: &Config{Level: "info"}}, + {label: "console tee", conf: &Config{Level: "debug"}, opts: withDiscardTee()}, + {label: "json tee", conf: &Config{JSON: true, Level: "debug"}, opts: withDiscardTee()}, + } + + for _, c := range cases { + b.Run(c.label, func(b *testing.B) { + root := c.logger(b) + + b.Run("plain", func(b *testing.B) { + b.ReportAllocs() + for range b.N { + root.Debugw("bench", "participant", "PA_1234") + } + }) + + b.Run("deferred", func(b *testing.B) { + l, resolver := root.WithDeferredValues() + resolver.Resolve("room", "RM_1234") + b.ReportAllocs() + for range b.N { + l.Debugw("bench", "participant", "PA_1234") + } + }) + + b.Run("deferred flush", func(b *testing.B) { + b.ReportAllocs() + for range b.N { + l, resolver := root.WithDeferredValues() + l.Debugw("bench", "participant", "PA_1234") + resolver.Resolve("room", "RM_1234") + } + }) + }) + } +} diff --git a/logger/logger_test.go b/logger/logger_test.go index 64393b800..60f30da4c 100644 --- a/logger/logger_test.go +++ b/logger/logger_test.go @@ -1,14 +1,16 @@ package logger import ( + "encoding/json" "fmt" + "os" "runtime" "strings" "testing" "github.com/stretchr/testify/require" - "go.uber.org/zap" "go.uber.org/zap/zapcore" + "go.uber.org/zap/zaptest/observer" "github.com/livekit/protocol/logger/testutil" "github.com/livekit/protocol/logger/zaputil" @@ -19,6 +21,22 @@ func zapLoggerCore(l Logger) zapcore.Core { return l.(ZapLogger).ToZap().Desugar().Core() } +func jsonTee(ws zapcore.WriteSyncer) ZapLoggerOption { + return WithTee(zaputil.NewTee(testutil.NewJSONCoreFactory(ws))) +} + +type componentLevels map[string]zapcore.Level + +func (c componentLevels) ResolveComponentLevel(component string) (zapcore.Level, bool) { + lvl, ok := c[component] + return lvl, ok +} + +func observerTee() (ZapLoggerOption, *observer.ObservedLogs) { + f, logs := testutil.NewObserverCoreFactory() + return WithTee(zaputil.NewTee(f)), logs +} + func TestLoggerComponent(t *testing.T) { t.Run("inheriting parent level", func(t *testing.T) { l, err := NewZapLogger(&Config{ @@ -83,8 +101,9 @@ func TestLoggerComponent(t *testing.T) { }) t.Run("log output matches expected values", func(t *testing.T) { + silenceStderr(t) ws := &testutil.BufferedWriteSyncer{} - l, err := NewZapLogger(&Config{}, WithTap(zaputil.NewWriteEnabler(ws, zapcore.DebugLevel))) + l, err := NewZapLogger(&Config{Level: "debug"}, jsonTee(ws)) require.NoError(t, err) l.Debugw("foo", "bar", "baz") @@ -98,21 +117,242 @@ func TestLoggerComponent(t *testing.T) { require.Equal(t, "baz", log.Bar) }) - t.Run("component enabler for tapped logger returns lowest enabled level", func(t *testing.T) { - tapLevel := zap.NewAtomicLevel() - l, err := NewZapLogger(&Config{Level: "info"}, WithTap(zaputil.NewWriteEnabler(&testutil.BufferedWriteSyncer{}, tapLevel))) + t.Run("component enabler ignores the tee", func(t *testing.T) { + tee, _ := observerTee() + l, err := NewZapLogger(&Config{Level: "info"}, tee) require.NoError(t, err) lvl := l.ComponentLeveler().ComponentLevel("foo") - // check config level require.False(t, lvl.Enabled(zapcore.DebugLevel)) require.True(t, lvl.Enabled(zapcore.InfoLevel)) + }) + + t.Run("branch leveler widens only the components it names", func(t *testing.T) { + l := must.Get(NewZapLogger(&Config{Level: "info"})) + lv := zaputil.NewComponentLeveler(l.Leveler(), componentLevels{"rtc.room": zapcore.DebugLevel}) + branch := l.WithComponentLeveler(lv) + + require.True(t, zapLoggerCore(branch.WithComponent("rtc").WithComponent("room")).Enabled(zapcore.DebugLevel)) + require.False(t, zapLoggerCore(branch.WithComponent("rtc")).Enabled(zapcore.DebugLevel)) + require.True(t, zapLoggerCore(branch.WithComponent("rtc")).Enabled(zapcore.InfoLevel)) + }) + + t.Run("branch leveler cannot quiet its parent", func(t *testing.T) { + l := must.Get(NewZapLogger(&Config{Level: "debug"})) + lv := zaputil.NewComponentLeveler(l.Leveler(), componentLevels{"sub": zapcore.ErrorLevel}) + + require.True(t, zapLoggerCore(l.WithComponentLeveler(lv).WithComponent("sub")).Enabled(zapcore.DebugLevel)) + }) + + t.Run("sibling branch levelers are independent", func(t *testing.T) { + l := must.Get(NewZapLogger(&Config{Level: "info"})) + a := zaputil.NewComponentLeveler(l.Leveler(), componentLevels{"sub": zapcore.DebugLevel}) + b := zaputil.NewComponentLeveler(l.Leveler(), componentLevels{}) + + require.True(t, zapLoggerCore(l.WithComponentLeveler(a).WithComponent("sub")).Enabled(zapcore.DebugLevel)) + require.False(t, zapLoggerCore(l.WithComponentLeveler(b).WithComponent("sub")).Enabled(zapcore.DebugLevel)) + }) + + t.Run("nil branch leveler is a no-op", func(t *testing.T) { + l := must.Get(NewZapLogger(&Config{Level: "info"})) + require.Same(t, Logger(l), l.WithComponentLeveler(nil)) + }) + + t.Run("write enablers are memoized per component", func(t *testing.T) { + lv := must.Get(NewZapLogger(&Config{Level: "info"})).Leveler() - // check tap level - tapLevel.SetLevel(zapcore.DebugLevel) - require.True(t, lvl.Enabled(zapcore.DebugLevel)) + require.Same(t, lv.WriteEnabler("sub"), lv.WriteEnabler("sub")) + require.NotSame(t, lv.WriteEnabler("sub"), lv.WriteEnabler("other")) }) + + t.Run("branch leveler follows global updates", func(t *testing.T) { + conf := &Config{Level: "info"} + l := must.Get(NewZapLogger(conf)) + lv := zaputil.NewComponentLeveler(l.Leveler(), componentLevels{"other": zapcore.DebugLevel}) + + core := zapLoggerCore(l.WithComponentLeveler(lv).WithComponent("sub")) + require.False(t, core.Enabled(zapcore.DebugLevel)) + + require.NoError(t, conf.Update(&Config{Level: "debug"})) + require.True(t, core.Enabled(zapcore.DebugLevel)) + }) + + t.Run("branch leveler refresh reaches existing loggers", func(t *testing.T) { + l := must.Get(NewZapLogger(&Config{Level: "info"})) + levels := componentLevels{} + lv := zaputil.NewComponentLeveler(l.Leveler(), levels) + + core := zapLoggerCore(l.WithComponentLeveler(lv).WithComponent("sub")) + require.False(t, core.Enabled(zapcore.DebugLevel)) + + levels["sub"] = zapcore.DebugLevel + lv.Refresh() + require.True(t, core.Enabled(zapcore.DebugLevel)) + + delete(levels, "sub") + lv.Refresh() + require.False(t, core.Enabled(zapcore.DebugLevel)) + }) +} + +func TestLoggerTee(t *testing.T) { + t.Run("receives WithValues fields per derived logger", func(t *testing.T) { + silenceStderr(t) + tee, logs := observerTee() + root := must.Get(NewZapLogger(&Config{Level: "debug"}, tee)) + + room := root.WithValues("room", "RM_1") + room.WithValues("participant", "PA_1").Debugw("participant") + room.WithValues("track", "TR_1").Debugw("track") + + entries := logs.All() + require.Len(t, entries, 2) + require.Equal(t, map[string]any{"room": "RM_1", "participant": "PA_1"}, entries[0].ContextMap()) + require.Equal(t, map[string]any{"room": "RM_1", "track": "TR_1"}, entries[1].ContextMap()) + }) + + t.Run("drops entries below the configured level", func(t *testing.T) { + tee, logs := observerTee() + l := must.Get(NewZapLogger(&Config{Level: "info"}, tee)) + + l.Debugw("debug") + + require.Empty(t, logs.All()) + }) + + t.Run("follows component levels", func(t *testing.T) { + silenceStderr(t) + tee, logs := observerTee() + l := must.Get(NewZapLogger(&Config{ + Level: "info", + ComponentLevels: map[string]string{ + "x": "debug", + "x.y": "info", + }, + }, tee)) + + x := l.WithComponent("x") + x.Debugw("x") + x.WithComponent("y").Debugw("xy") + + entries := logs.All() + require.Len(t, entries, 1) + require.Equal(t, "x", entries[0].Message) + }) + + t.Run("follows the branch leveler", func(t *testing.T) { + silenceStderr(t) + tee, logs := observerTee() + l := must.Get(NewZapLogger(&Config{Level: "info"}, tee)) + lv := zaputil.NewComponentLeveler(l.Leveler(), FixedComponentLevel(zapcore.DebugLevel)) + + l.Debugw("dropped") + l.WithComponentLeveler(lv).Debugw("kept") + + entries := logs.All() + require.Len(t, entries, 1) + require.Equal(t, "kept", entries[0].Message) + }) + + t.Run("receives resolved deferred values", func(t *testing.T) { + silenceStderr(t) + tee, logs := observerTee() + root := must.Get(NewZapLogger(&Config{Level: "debug"}, tee)) + + l, resolver := root.WithDeferredValues() + l.Debugw("deferred") + require.Empty(t, logs.All()) + + resolver.Resolve("participant", "PA_1") + + entries := logs.All() + require.Len(t, entries, 1) + require.Equal(t, map[string]any{"participant": "PA_1"}, entries[0].ContextMap()) + }) + + t.Run("deferred entries honor the resolved level", func(t *testing.T) { + silenceStderr(t) + tee, logs := observerTee() + root := must.Get(NewZapLogger(&Config{Level: "warn"}, tee)) + + l, resolver := root.WithDeferredValues() + l.Debugw("below") + l.Warnw("at", nil) + resolver.Resolve("participant", "PA_1") + + entries := logs.All() + require.Len(t, entries, 1) + require.Equal(t, "at", entries[0].Message) + }) + + t.Run("deferred entries run the tee's check-time logic", func(t *testing.T) { + silenceStderr(t) + core, logs := observer.New(zapcore.DebugLevel) + var hooked int + tee := zaputil.NewTee(func(enab zapcore.LevelEnabler) zapcore.Core { + return zapcore.RegisterHooks(testutil.Leveled(core, enab), func(zapcore.Entry) error { + hooked++ + return nil + }) + }) + root := must.Get(NewZapLogger(&Config{Level: "debug"}, WithTee(tee))) + + l, resolver := root.WithDeferredValues() + l.Debugw("deferred") + resolver.Resolve("participant", "PA_1") + + require.Len(t, logs.All(), 1) + require.Equal(t, 1, hooked) + }) + + t.Run("agrees with the console on malformed value lists", func(t *testing.T) { + readStderr := captureStderr(t) + tee, logs := observerTee() + l := must.Get(NewZapLogger(&Config{JSON: true, Level: "debug"}, tee)) + + l.WithValues("room", "RM_1", 42, "dropped", "track", "TR_1", "dangling").Debugw("test") + + var console map[string]any + require.NoError(t, json.Unmarshal([]byte(readStderr()), &console)) + for _, k := range []string{"level", "ts", "caller", "msg"} { + delete(console, k) + } + + require.Equal(t, map[string]any{"room": "RM_1", "track": "TR_1"}, console) + require.Equal(t, console, logs.All()[0].ContextMap()) + }) +} + +// The core captures os.Stderr when it is built, so these must run before the logger is created. + +func silenceStderr(tb testing.TB) { + tb.Helper() + f, err := os.OpenFile(os.DevNull, os.O_WRONLY, 0) + require.NoError(tb, err) + prev := os.Stderr + os.Stderr = f + tb.Cleanup(func() { + os.Stderr = prev + f.Close() + }) +} + +func captureStderr(tb testing.TB) func() string { + tb.Helper() + f, err := os.CreateTemp(tb.TempDir(), "stderr") + require.NoError(tb, err) + prev := os.Stderr + os.Stderr = f + tb.Cleanup(func() { + os.Stderr = prev + f.Close() + }) + return func() string { + b, err := os.ReadFile(f.Name()) + require.NoError(tb, err) + return string(b) + } } type TestLogOutput struct { @@ -160,7 +400,7 @@ func TestLoggerCallDepth(t *testing.T) { for label, getLogFunc := range cases { t.Run(label, func(t *testing.T) { ws := &testutil.BufferedWriteSyncer{} - l := must.Get(NewZapLogger(&Config{}, WithTap(zaputil.NewWriteEnabler(ws, zapcore.DebugLevel)))) + l := must.Get(NewZapLogger(&Config{Level: "debug"}, jsonTee(ws))) testLogCaller(getLogFunc(l)) diff --git a/logger/testutil/testutil.go b/logger/testutil/testutil.go index 7b7a11459..379f2baf8 100644 --- a/logger/testutil/testutil.go +++ b/logger/testutil/testutil.go @@ -3,6 +3,10 @@ package testutil import ( "bytes" "encoding/json" + + "go.uber.org/zap" + "go.uber.org/zap/zapcore" + "go.uber.org/zap/zaptest/observer" ) type TestLogOutput struct { @@ -21,3 +25,50 @@ func (t *BufferedWriteSyncer) Unmarshal(v any) error { } func (t *BufferedWriteSyncer) Sync() error { return nil } + +func NewJSONCore(ws zapcore.WriteSyncer, enab zapcore.LevelEnabler) zapcore.Core { + return zapcore.NewCore(zapcore.NewJSONEncoder(zap.NewProductionEncoderConfig()), ws, enab) +} + +// zaputil's own tests import this package, so it must not import zaputil; +// callers build the Tee themselves. +type CoreFactory = func(enab zapcore.LevelEnabler) zapcore.Core + +func NewJSONCoreFactory(ws zapcore.WriteSyncer) CoreFactory { + return func(enab zapcore.LevelEnabler) zapcore.Core { + return NewJSONCore(ws, enab) + } +} + +// NewObserverCoreFactory records every entry the core is handed. +func NewObserverCoreFactory() (CoreFactory, *observer.ObservedLogs) { + core, logs := observer.New(zapcore.DebugLevel) + return func(enab zapcore.LevelEnabler) zapcore.Core { + return Leveled(core, enab) + }, logs +} + +// Leveled overrides a fixed-level core's level decisions with enab. +func Leveled(core zapcore.Core, enab zapcore.LevelEnabler) zapcore.Core { + return leveledCore{core, enab} +} + +type leveledCore struct { + zapcore.Core + enab zapcore.LevelEnabler +} + +func (c leveledCore) Enabled(lvl zapcore.Level) bool { + return c.enab.Enabled(lvl) +} + +func (c leveledCore) With(fields []zapcore.Field) zapcore.Core { + return leveledCore{c.Core.With(fields), c.enab} +} + +func (c leveledCore) Check(ent zapcore.Entry, ce *zapcore.CheckedEntry) *zapcore.CheckedEntry { + if !c.enab.Enabled(ent.Level) { + return ce + } + return ce.AddCore(ent, c) +} diff --git a/logger/zaputil/deferrer.go b/logger/zaputil/deferrer.go index dad28b0ae..09e07d03d 100644 --- a/logger/zaputil/deferrer.go +++ b/logger/zaputil/deferrer.go @@ -62,18 +62,22 @@ func (b *Deferrer) flush() { n := len(fields) for _, w := range writes { - fields = append(fields[:n], w.fields...) - w.core.Write(w.ent, fields) + if ce := w.core.Check(w.ent, nil); ce != nil { + ce.Write(append(fields[:n], w.fields...)...) + } } } -func (b *Deferrer) write(core zapcore.Core, ent zapcore.Entry, fields []zapcore.Field) error { +func (b *Deferrer) write(core zapcore.Core, ent zapcore.Entry, fields []zapcore.Field) { for { if dfs := b.fields.Load(); dfs != nil { - return core.Write(ent, slices.Concat(fields, *dfs)) + if ce := core.Check(ent, nil); ce != nil { + ce.Write(slices.Concat(fields, *dfs)...) + } + return } if b.buffer(core, ent, fields) { - return nil + return } } } @@ -158,5 +162,6 @@ func (c *deferredValueCore) Check(ent zapcore.Entry, ce *zapcore.CheckedEntry) * } func (c *deferredValueCore) Write(ent zapcore.Entry, fields []zapcore.Field) error { - return c.def.write(c.Core, ent, fields) + c.def.write(c.Core, ent, fields) + return nil } diff --git a/logger/zaputil/deferrer_test.go b/logger/zaputil/deferrer_test.go index 6b6af9b1d..65e73dfa6 100644 --- a/logger/zaputil/deferrer_test.go +++ b/logger/zaputil/deferrer_test.go @@ -15,6 +15,7 @@ package zaputil import ( + "io" "testing" "github.com/stretchr/testify/require" @@ -50,9 +51,8 @@ func TestDeferredLogger(t *testing.T) { t.Run("resolved values can be overwritten", func(t *testing.T) { ws := &testutil.BufferedWriteSyncer{} - we := NewWriteEnabler(ws, zapcore.DebugLevel) enc := zapcore.NewJSONEncoder(zap.NewProductionEncoderConfig()) - c := NewEncoderCore(enc, we) + c := zapcore.NewCore(enc, ws, zapcore.DebugLevel) d := &Deferrer{} dc := NewDeferredValueCore(c, d) s := zap.New(dc).Sugar() @@ -77,9 +77,8 @@ func TestDeferredLogger(t *testing.T) { t.Run("resolved values merge with previous resolutions", func(t *testing.T) { ws := &testutil.BufferedWriteSyncer{} - we := NewWriteEnabler(ws, zapcore.DebugLevel) enc := zapcore.NewJSONEncoder(zap.NewProductionEncoderConfig()) - c := NewEncoderCore(enc, we) + c := zapcore.NewCore(enc, ws, zapcore.DebugLevel) d := &Deferrer{} dc := NewDeferredValueCore(c, d) s := zap.New(dc).Sugar() @@ -99,9 +98,8 @@ func TestDeferredLogger(t *testing.T) { t.Run("re-resolve", func(t *testing.T) { ws := &testutil.BufferedWriteSyncer{} - we := NewWriteEnabler(ws, zapcore.DebugLevel) enc := zapcore.NewJSONEncoder(zap.NewProductionEncoderConfig()) - c := NewEncoderCore(enc, we) + c := zapcore.NewCore(enc, ws, zapcore.DebugLevel) d := &Deferrer{} dc := NewDeferredValueCore(c, d) s := zap.New(dc).Sugar() @@ -130,4 +128,32 @@ func TestDeferredLogger(t *testing.T) { require.Equal(t, "car", log.A) require.Equal(t, "dog", log.B) }) + + t.Run("each destination applies its own level", func(t *testing.T) { + debug := countingCore(zapcore.DebugLevel) + warn := countingCore(zapcore.WarnLevel) + d := &Deferrer{} + s := zap.New(NewDeferredValueCore(zapcore.NewTee(debug, warn), d)).Sugar() + + s.Infow("test") + d.Resolve("a", "foo") + require.Equal(t, 1, debug.WriteCount()) + require.Equal(t, 0, warn.WriteCount()) + + s.Infow("test") + require.Equal(t, 2, debug.WriteCount()) + require.Equal(t, 0, warn.WriteCount()) + + s.Warnw("test") + require.Equal(t, 3, debug.WriteCount()) + require.Equal(t, 1, warn.WriteCount()) + }) +} + +func countingCore(enab zapcore.LevelEnabler) *testCore { + return &testCore{Core: zapcore.NewCore( + zapcore.NewJSONEncoder(zap.NewProductionEncoderConfig()), + zapcore.AddSync(io.Discard), + enab, + )} } diff --git a/logger/zaputil/encoder.go b/logger/zaputil/encoder.go index 909645084..095c386e3 100644 --- a/logger/zaputil/encoder.go +++ b/logger/zaputil/encoder.go @@ -14,10 +14,7 @@ package zaputil -import ( - "go.uber.org/multierr" - "go.uber.org/zap/zapcore" -) +import "go.uber.org/zap/zapcore" type WriteEnabler struct { zapcore.WriteSyncer @@ -27,88 +24,3 @@ type WriteEnabler struct { func NewWriteEnabler(ws zapcore.WriteSyncer, enab zapcore.LevelEnabler) *WriteEnabler { return &WriteEnabler{ws, enab} } - -type discardWriteSyncer struct{} - -func (discardWriteSyncer) Write(p []byte) (int, error) { return len(p), nil } -func (discardWriteSyncer) Sync() error { return nil } - -func NewDiscardWriteEnabler() *WriteEnabler { - return NewWriteEnabler(discardWriteSyncer{}, zapcore.FatalLevel) -} - -func NewEncoderCore(enc zapcore.Encoder, out ...*WriteEnabler) zapcore.Core { - return &encoderCore{ - enc: enc, - out: out, - } -} - -type encoderCore struct { - enc zapcore.Encoder - out []*WriteEnabler -} - -func (c encoderCore) Level() zapcore.Level { - minLvl := zapcore.FatalLevel - for _, out := range c.out { - if lvl := zapcore.LevelOf(out); lvl < minLvl { - minLvl = lvl - } - } - return minLvl -} - -func (c encoderCore) Enabled(lvl zapcore.Level) bool { - for _, out := range c.out { - if out.Enabled(lvl) { - return true - } - } - return false -} - -func (c *encoderCore) With(fields []zapcore.Field) zapcore.Core { - dup := *c - dup.enc = dup.enc.Clone() - for _, f := range fields { - f.AddTo(dup.enc) - } - return &dup -} - -func (c *encoderCore) Check(ent zapcore.Entry, ce *zapcore.CheckedEntry) *zapcore.CheckedEntry { - if c.Enabled(ent.Level) { - return ce.AddCore(ent, c) - } - return ce -} - -func (c *encoderCore) Write(ent zapcore.Entry, fields []zapcore.Field) error { - buf, err := c.enc.EncodeEntry(ent, fields) - if err != nil { - return err - } - for _, out := range c.out { - if out.Enabled(ent.Level) { - _, werr := out.Write(buf.Bytes()) - err = multierr.Append(err, werr) - } - } - buf.Free() - if err != nil { - return err - } - if ent.Level> zapcore.ErrorLevel { - _ = c.Sync() - } - return nil -} - -func (c *encoderCore) Sync() error { - var err error - for _, out := range c.out { - err = multierr.Append(err, out.Sync()) - } - return err -} diff --git a/logger/zaputil/leveler.go b/logger/zaputil/leveler.go new file mode 100644 index 000000000..4a5959847 --- /dev/null +++ b/logger/zaputil/leveler.go @@ -0,0 +1,121 @@ +// Copyright 2023 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package zaputil + +import ( + "sync" + + "github.com/puzpuzpuz/xsync/v4" + "go.uber.org/zap" + "go.uber.org/zap/zapcore" +) + +// ComponentLevelResolver is given a full dotted component path. Returning false defers +// to the leveler's parent. +type ComponentLevelResolver interface { + ResolveComponentLevel(component string) (zapcore.Level, bool) +} + +type FixedComponentLevel zapcore.Level + +func (l FixedComponentLevel) ResolveComponentLevel(string) (zapcore.Level, bool) { + return zapcore.Level(l), true +} + +// level is held apart from enab because only the former is settable once handed out. +type componentEntry struct { + level zap.AtomicLevel + enab zapcore.LevelEnabler + write *WriteEnabler +} + +// One ComponentLeveler per configuration source. +type ComponentLeveler struct { + parent *ComponentLeveler + resolver ComponentLevelResolver + ws zapcore.WriteSyncer + + mu sync.Mutex + entries *xsync.Map[string, *componentEntry] +} + +func NewRootComponentLeveler(ws zapcore.WriteSyncer, r ComponentLevelResolver) *ComponentLeveler { + return &ComponentLeveler{ + resolver: r, + ws: ws, + entries: xsync.NewMap[string, *componentEntry](), + } +} + +// NewComponentLeveler derives a leveler that can only widen parent: a level parent +// already enables stays enabled whatever r resolves. +func NewComponentLeveler(parent *ComponentLeveler, r ComponentLevelResolver) *ComponentLeveler { + return &ComponentLeveler{ + parent: parent, + resolver: r, + ws: parent.ws, + entries: xsync.NewMap[string, *componentEntry](), + } +} + +func (l *ComponentLeveler) ComponentLevel(component string) zapcore.LevelEnabler { + return l.entry(component).enab +} + +// WriteEnabler is memoized per component so that rebuilding a core does not allocate. +func (l *ComponentLeveler) WriteEnabler(component string) *WriteEnabler { + return l.entry(component).write +} + +// Refresh re-resolves in place, so enablers already handed out see the new levels. +func (l *ComponentLeveler) Refresh() { + l.mu.Lock() + defer l.mu.Unlock() + l.entries.Range(func(component string, e *componentEntry) bool { + e.level.SetLevel(l.resolve(component)) + return true + }) +} + +func (l *ComponentLeveler) entry(component string) *componentEntry { + if e, ok := l.entries.Load(component); ok { + return e + } + + // Serialized against Refresh, otherwise an entry resolved here could miss a + // concurrent configuration change and stay stale forever. + l.mu.Lock() + defer l.mu.Unlock() + if e, ok := l.entries.Load(component); ok { + return e + } + + e := &componentEntry{level: zap.NewAtomicLevelAt(l.resolve(component))} + e.enab = e.level + if l.parent != nil { + e.enab = OrLevelEnabler{l.parent.ComponentLevel(component), e.level} + } + e.write = NewWriteEnabler(l.ws, e.enab) + l.entries.Store(component, e) + return e +} + +// InvalidLevel enables nothing, leaving the Or against the parent decisive. +func (l *ComponentLeveler) resolve(component string) zapcore.Level { + if lvl, ok := l.resolver.ResolveComponentLevel(component); ok { + return lvl + } + return zapcore.InvalidLevel +} diff --git a/logger/zaputil/leveler_test.go b/logger/zaputil/leveler_test.go new file mode 100644 index 000000000..05b154249 --- /dev/null +++ b/logger/zaputil/leveler_test.go @@ -0,0 +1,96 @@ +package zaputil + +import ( + "fmt" + "io" + "sync" + "testing" + + "github.com/stretchr/testify/require" + "go.uber.org/zap/zapcore" +) + +type componentLevels map[string]zapcore.Level + +func (c componentLevels) ResolveComponentLevel(component string) (zapcore.Level, bool) { + lvl, ok := c[component] + return lvl, ok +} + +func discardRoot(lvl zapcore.Level) *ComponentLeveler { + return NewRootComponentLeveler(zapcore.AddSync(io.Discard), FixedComponentLevel(lvl)) +} + +func TestComponentLevelerMemoizes(t *testing.T) { + root := discardRoot(zapcore.InfoLevel) + + require.Same(t, root.WriteEnabler("sub"), root.WriteEnabler("sub")) + require.NotSame(t, root.WriteEnabler("sub"), root.WriteEnabler("other")) + require.Equal(t, root.ComponentLevel("sub"), root.ComponentLevel("sub")) +} + +func TestComponentLevelerWidensParent(t *testing.T) { + root := discardRoot(zapcore.InfoLevel) + child := NewComponentLeveler(root, componentLevels{"sub": zapcore.DebugLevel}) + + require.True(t, child.ComponentLevel("sub").Enabled(zapcore.DebugLevel)) + require.False(t, child.ComponentLevel("other").Enabled(zapcore.DebugLevel)) + require.True(t, child.ComponentLevel("other").Enabled(zapcore.InfoLevel)) +} + +func TestComponentLevelerCannotQuietParent(t *testing.T) { + root := discardRoot(zapcore.DebugLevel) + child := NewComponentLeveler(root, componentLevels{"sub": zapcore.ErrorLevel}) + + require.True(t, child.ComponentLevel("sub").Enabled(zapcore.DebugLevel)) +} + +func TestComponentLevelerRefresh(t *testing.T) { + root := discardRoot(zapcore.InfoLevel) + levels := componentLevels{} + child := NewComponentLeveler(root, levels) + + enab := child.ComponentLevel("sub") + require.False(t, enab.Enabled(zapcore.DebugLevel)) + + levels["sub"] = zapcore.DebugLevel + child.Refresh() + require.True(t, enab.Enabled(zapcore.DebugLevel)) + + delete(levels, "sub") + child.Refresh() + require.False(t, enab.Enabled(zapcore.DebugLevel)) + require.True(t, enab.Enabled(zapcore.InfoLevel)) +} + +// Refresh walks the entry map while other goroutines materialize into it. +func TestComponentLevelerConcurrentRefresh(t *testing.T) { + root := discardRoot(zapcore.InfoLevel) + child := NewComponentLeveler(root, FixedComponentLevel(zapcore.DebugLevel)) + + var wg sync.WaitGroup + for i := range 8 { + wg.Add(1) + go func() { + defer wg.Done() + for j := range 200 { + c := fmt.Sprintf("c%d.s%d", i, j%16) + if child.WriteEnabler(c) == nil || root.ComponentLevel(c) == nil { + t.Error("nil enabler") + return + } + } + }() + } + + wg.Add(1) + go func() { + defer wg.Done() + for range 200 { + root.Refresh() + child.Refresh() + } + }() + + wg.Wait() +} diff --git a/logger/zaputil/tee.go b/logger/zaputil/tee.go new file mode 100644 index 000000000..e7325b32a --- /dev/null +++ b/logger/zaputil/tee.go @@ -0,0 +1,55 @@ +// Copyright 2023 LiveKit, Inc. +// +// Licensed under the Apache License, Version 2.0 (the "License"); +// you may not use this file except in compliance with the License. +// You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, +// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +// See the License for the specific language governing permissions and +// limitations under the License. + +package zaputil + +import ( + "slices" + + "go.uber.org/zap/zapcore" +) + +// Tee duplicates every entry a derived logger writes. The enabler handed to +// newCore is the level that logger resolved, and is what the built core must +// gate on. +type Tee struct { + newCore func(enab zapcore.LevelEnabler) zapcore.Core + fields []zapcore.Field +} + +func NewTee(newCore func(enab zapcore.LevelEnabler) zapcore.Core) Tee { + return Tee{newCore: newCore} +} + +// Core returns nil for the zero Tee. +func (t Tee) Core(enab zapcore.LevelEnabler) zapcore.Core { + if t.newCore == nil { + return nil + } + core := t.newCore(enab) + if len(t.fields) == 0 { + return core + } + return core.With(t.fields) +} + +func (t Tee) WithValues(kvs ...any) Tee { + if t.newCore == nil { + return t + } + // Clip so sibling loggers derived from this one cannot write into a shared + // backing array. + t.fields = append(slices.Clip(t.fields), valueFields(kvs...)...) + return t +} diff --git a/logger/zaputil/zaputil.go b/logger/zaputil/zaputil.go index bf7b36332..b6f5cf638 100644 --- a/logger/zaputil/zaputil.go +++ b/logger/zaputil/zaputil.go @@ -29,48 +29,34 @@ func encoderWithValues(enc zapcore.Encoder, kvs ...any) zapcore.Encoder { return clone } -type Encoder[T any] interface { - WithValues(kvs ...any) T - Core(console, json *WriteEnabler) zapcore.Core -} - -type DevelopmentEncoder struct { - console zapcore.Encoder - json zapcore.Encoder -} - -func NewDevelopmentEncoder() DevelopmentEncoder { - return DevelopmentEncoder{ - console: zapcore.NewConsoleEncoder(zap.NewDevelopmentEncoderConfig()), - json: zapcore.NewJSONEncoder(zap.NewProductionEncoderConfig()), +// Field selection must match encoderWithValues, which cannot share this loop without allocating. +func valueFields(kvs ...any) []zapcore.Field { + fields := make([]zapcore.Field, 0, len(kvs)/2) + for i := 1; i < len(kvs); i += 2 { + if key, ok := kvs[i-1].(string); ok { + fields = append(fields, zap.Any(key, kvs[i])) + } } + return fields } -func (e DevelopmentEncoder) WithValues(kvs ...any) DevelopmentEncoder { - e.console = encoderWithValues(e.console, kvs...) - e.json = encoderWithValues(e.json, kvs...) - return e +type Encoder struct { + enc zapcore.Encoder } -func (e DevelopmentEncoder) Core(console, json *WriteEnabler) zapcore.Core { - return zapcore.NewTee(NewEncoderCore(e.console, console), NewEncoderCore(e.json, json)) +func NewDevelopmentEncoder() Encoder { + return Encoder{zapcore.NewConsoleEncoder(zap.NewDevelopmentEncoderConfig())} } -type ProductionEncoder struct { - json zapcore.Encoder -} - -func NewProductionEncoder() ProductionEncoder { - return ProductionEncoder{ - json: zapcore.NewJSONEncoder(zap.NewProductionEncoderConfig()), - } +func NewProductionEncoder() Encoder { + return Encoder{zapcore.NewJSONEncoder(zap.NewProductionEncoderConfig())} } -func (e ProductionEncoder) WithValues(kvs ...any) ProductionEncoder { - e.json = encoderWithValues(e.json, kvs...) +func (e Encoder) WithValues(kvs ...any) Encoder { + e.enc = encoderWithValues(e.enc, kvs...) return e } -func (e ProductionEncoder) Core(console, json *WriteEnabler) zapcore.Core { - return NewEncoderCore(e.json, console, json) +func (e Encoder) Core(out *WriteEnabler) zapcore.Core { + return zapcore.NewCore(e.enc, out, out) } diff --git a/logger/zaputil/zaputil_test.go b/logger/zaputil/zaputil_test.go index bbbb6c511..1c3c619bb 100644 --- a/logger/zaputil/zaputil_test.go +++ b/logger/zaputil/zaputil_test.go @@ -51,6 +51,15 @@ func (c *testCore) With(fields []zapcore.Field) zapcore.Core { } } +// Check must register c, not the embedded core, or Write never runs +func (c *testCore) Check(ent zapcore.Entry, ce *zapcore.CheckedEntry) *zapcore.CheckedEntry { + c.init() + if c.Enabled(ent.Level) { + return ce.AddCore(ent, c) + } + return ce +} + func (c *testCore) Write(entry zapcore.Entry, fields []zapcore.Field) error { c.init() c.writeCount.Inc()

AltStyle によって変換されたページ (->オリジナル) /