From 271d9cde38974b023239c60c7f7b04f1b1a33e6a Mon Sep 17 00:00:00 2001 From: Paul Wells Date: Thu, 10 Sep 2026 05:12:19 -0700 Subject: [PATCH] logger: replace WithTap with WithTee (#1790) * logger/zaputil: replay deferred writes through Check Deferrer.write and Deferrer.flush called the wrapped core's Write directly, bypassing Check. zapcore.multiCore.Write fans out to every child unconditionally, so a destination that filters in Check rather than in Write received entries it never enabled. encoderCore was immune because it re-checks each out's level inside Write; nothing else is, including zaptest/observer. testCore now registers itself in Check. It overrides Write but not Check, so its Write only ever ran because the deferred path skipped Check. * logger: replace WithTap with WithTee WithTap took a WriteSyncer, so a tap could only ever receive bytes encoded by a zaputil-chosen encoder: Encoder.Core took two destinations and DevelopmentEncoder carried a second JSON encoder that existed only to format tap output. WithTee hands the duplicate stream a zapcore.Core instead, which brings its own encoder and sees the structured entry. A Tee builds its core from the derived logger's resolved enabler, so it shares the console's level rather than widening what the logger emits, and the component leveler no longer has to OR a tap's level in. Fields do not arrive on their own: WithValues bakes them into the encoder and replaces the whole core, and the sugared path never calls Core.With, so a Tee accumulates them itself. Benchmarks cover objects allocated per WithValues and WithDeferredValues, with and without a tee. * logger/zaputil: replace encoderCore with zapcore.NewCore With one destination per core, encoderCore is zapcore.ioCore plus a per-destination level re-check on write, and that check became redundant once the deferred path started replaying through Check. * logger/zaputil: drop the Encoder type parameter Encoder was generic because WithValues returned a concrete type, and there were two concrete types because DevelopmentEncoder needed a second encoder for the tap. Both are now one encoder and one destination, so the interface and its two implementations collapse into a struct with two constructors, and zapLogger, newZapLogger and zapLoggerComponentLeveler stop being generic. * logger: make the component leveler replaceable per branch Component levels resolved through one process-global sharedConfig, so the only per-tenant knob was WithMinLevel: a flat floor that also bypassed the write-enabler cache and allocated a WriteEnabler on every derivation. zaputil.ComponentLeveler now owns both the per-component level and the write-enabler cache, and WithComponentLeveler attaches one to a branch of the logger tree. A derived leveler can only widen its parent, which is what lets a global config change reach every branch without the root holding references to its children. Replaces WithMinLevel. * changeset: cover WithComponentLeveler * logger: tighten the tee comments --- .changeset/logger-with-tee.md | 8 + livekit/egress_test.go | 3 +- logger/config.go | 22 ++- logger/config_test.go | 25 +++ logger/logger.go | 194 +++++++++--------------- logger/logger_bench_test.go | 166 ++++++++++++++++++++ logger/logger_test.go | 260 ++++++++++++++++++++++++++++++-- logger/testutil/testutil.go | 51 +++++++ logger/zaputil/deferrer.go | 17 ++- logger/zaputil/deferrer_test.go | 38 ++++- logger/zaputil/encoder.go | 90 +---------- logger/zaputil/leveler.go | 121 +++++++++++++++ logger/zaputil/leveler_test.go | 96 ++++++++++++ logger/zaputil/tee.go | 55 +++++++ logger/zaputil/zaputil.go | 50 +++--- logger/zaputil/zaputil_test.go | 9 ++ 16 files changed, 933 insertions(+), 272 deletions(-) create mode 100644 .changeset/logger-with-tee.md create mode 100644 logger/config_test.go create mode 100644 logger/logger_bench_test.go create mode 100644 logger/zaputil/leveler.go create mode 100644 logger/zaputil/leveler_test.go create mode 100644 logger/zaputil/tee.go 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()