From 7948c85b81a818e1059794e612569c16b0a85895 Mon Sep 17 00:00:00 2001 From: Cosmin Staicu Date: Wed, 23 Sep 2026 21:19:50 +0300 Subject: [PATCH] feat(logger): apply a log config file to the live logger without a restart MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The pieces for changing levels at runtime were already here — every component level is a zap.AtomicLevel held by the root ComponentLeveler, and FromZapLogger registers a Config observer that refreshes the leveler — but nothing ever changed a level on a running process, so a level change meant restarting it. Add the missing trigger: when LK_LOG_CONFIG_PATH is set, poll that file and apply its level and component_levels to the Config the service is holding (LK_LOG_CONFIG_INTERVAL overrides the 30s default). The hook sits in FromZapLogger, which is the one path every consumer reaches — livekit-server via InitFromConfig, livekit-sip via NewZapLogger — so no binary needs its own flag or call site. Polling rather than fsnotify because the target is a mounted ConfigMap: kubelet swaps the ..data symlink instead of rewriting the file, so a watch on the file never fires. Only those two keys are taken. The rest of the logging block (json, sample, the sampler settings) is read once when a logger is built, and WithItemSampler reads the item sampler settings without Config's lock, so rewriting them on a running process would change nothing and race that read. The file is decoded strictly, so one that names any other key is rejected with a warning rather than half applied, and the levels are written through Config.updateLevels, which touches nothing else. The read side needs nothing new: the leveler resolves through Config.ResolveComponentLevel, which takes the same lock. The file is a declarative overlay on the startup config, not on whatever was applied last: WatchConfigFile snapshots the config once at startup and every file is applied over that baseline, so a level the file omits falls back to its startup value and emptying the file to `{}` restores the levels the process booted with. component_levels entries merge over the startup ones, so an entry the file does not name — including pion_level, which livekit-server puts there — keeps its startup value, and one that disappears from the file stops applying. An unreadable file (an optional ConfigMap not yet mounted), unchanged bytes and a rejected file all leave the config in force untouched, and a rejected file is reported once when it appears rather than on every poll. A file that goes away is deliberately not a reset — a transient read error would otherwise flap levels on a live process; emptying it to `{}` is the reset. Signed-off-by: Cosmin Staicu --- .changeset/logger-config-file-watch.md | 6 + logger/config.go | 42 +++++ logger/configwatch.go | 141 ++++++++++++++++ logger/configwatch_test.go | 219 +++++++++++++++++++++++++ logger/logger.go | 1 + 5 files changed, 409 insertions(+) create mode 100644 .changeset/logger-config-file-watch.md create mode 100644 logger/configwatch.go create mode 100644 logger/configwatch_test.go diff --git a/.changeset/logger-config-file-watch.md b/.changeset/logger-config-file-watch.md new file mode 100644 index 000000000..28e56e132 --- /dev/null +++ b/.changeset/logger-config-file-watch.md @@ -0,0 +1,6 @@ +--- +"github.com/livekit/protocol": patch +"@livekit/protocol": patch +--- + +logger: apply a log config file to the live logger without a restart, by setting `LK_LOG_CONFIG_PATH`. diff --git a/logger/config.go b/logger/config.go index 75c77477e..1bef59086 100644 --- a/logger/config.go +++ b/logger/config.go @@ -72,6 +72,48 @@ func (c *Config) Update(o *Config) error { return nil } +// updateLevels replaces only the levels and notifies the same observers as Update. The other +// fields are read once when a logger is built, and WithItemSampler reads the item sampler +// settings without the lock, so a running process must never have them rewritten. +func (c *Config) updateLevels(level string, componentLevels map[string]string) error { + c.lock.Lock() + c.Level = level + c.ComponentLevels = componentLevels + callbacks := c.onUpdatedCallbacks + c.lock.Unlock() + + for _, cb := range callbacks { + if err := cb(c); err != nil { + return err + } + } + return nil +} + +// snapshot copies the data fields so a file can be unmarshalled over the config in force. +// Update assigns every field, so decoding a partial file into a zero Config would silently reset +// the rest — including ComponentLevels, which is where livekit-server puts pion_level. +func (c *Config) snapshot() *Config { + c.lock.Lock() + defer c.lock.Unlock() + + componentLevels := make(map[string]string, len(c.ComponentLevels)) + for component, level := range c.ComponentLevels { + componentLevels[component] = level + } + return &Config{ + JSON: c.JSON, + Level: c.Level, + Sample: c.Sample, + ComponentLevels: componentLevels, + SampleInitial: c.SampleInitial, + SampleInterval: c.SampleInterval, + ItemSampleSeconds: c.ItemSampleSeconds, + ItemSampleInitial: c.ItemSampleInitial, + ItemSampleInterval: c.ItemSampleInterval, + } +} + func (c *Config) AddUpdateObserver(cb ConfigObserver) { c.lock.Lock() defer c.lock.Unlock() diff --git a/logger/configwatch.go b/logger/configwatch.go new file mode 100644 index 000000000..0114b3084 --- /dev/null +++ b/logger/configwatch.go @@ -0,0 +1,141 @@ +package logger + +import ( + "bytes" + "errors" + "io" + "maps" + "os" + "sync" + "time" + + "gopkg.in/yaml.v3" +) + +const ( + // ConfigPathEnv names a file holding the `level` and `component_levels` keys of the service + // config's `logging` block. Point it inside a mounted ConfigMap to change levels without + // restarting: kubelet refreshes the mount in place, and the next poll pushes the new values + // into the live logger. + ConfigPathEnv = "LK_LOG_CONFIG_PATH" + // ConfigIntervalEnv overrides the poll interval as a Go duration (e.g. "10s"). + ConfigIntervalEnv = "LK_LOG_CONFIG_INTERVAL" + + defaultConfigWatchInterval = 30 * time.Second +) + +var configWatchOnce sync.Once + +// startConfigWatchFromEnv wires the watcher for the first logger the process builds, which is the +// one whose Config the service keeps. Every binary that uses this package reaches it through +// FromZapLogger, so none of them need their own flag or call site. +func startConfigWatchFromEnv(conf *Config) { + path := os.Getenv(ConfigPathEnv) + if path == "" { + return + } + interval := defaultConfigWatchInterval + if v := os.Getenv(ConfigIntervalEnv); v != "" { + if d, err := time.ParseDuration(v); err == nil && d > 0 { + interval = d + } + } + configWatchOnce.Do(func() { + WatchConfigFile(conf, path, interval) + }) +} + +// WatchConfigFile applies path to conf every interval until the returned stop is called. +// +// Polling rather than fsnotify on purpose: a ConfigMap volume update swaps the `..data` symlink +// instead of rewriting the file, so a watch on the file itself never fires. +// The returned stop is synchronous: once it returns, no further apply can be in flight. +func WatchConfigFile(conf *Config, path string, interval time.Duration) (stop func()) { + done := make(chan struct{}) + stopped := make(chan struct{}) + // The config the process started with. Every file is applied over this, never over whatever + // the previous file left in force, so an empty file restores the startup levels and a + // component_levels entry that disappears from the file stops applying. + baseline := conf.snapshot() + go func() { + defer close(stopped) + ticker := time.NewTicker(interval) + defer ticker.Stop() + var last []byte + for { + select { + case <-done: + return + case <-ticker.C: + // A tick and a close can be ready at once and select would pick either, + // so re-check: after stop the file must not be applied again. + select { + case <-done: + return + default: + } + if read, _ := applyConfigFile(conf, baseline, path, last); read != nil { + last = read + } + } + } + }() + + var once sync.Once + return func() { + once.Do(func() { close(done) }) + <-stopped + } +} + +// applyConfigFile overlays the levels in path on baseline and pushes them into conf when the bytes +// differ from last, the ones the previous call read. It returns the bytes it read, nil when the file could +// not be read, and whether it applied them. ok is false when nothing was applied (unreadable, +// unchanged or invalid), and the config already in force stays untouched. +// +// Invalid bytes come back as read too, so a caller that passes them in next time skips them: a bad +// file is reported once when it appears, not on every poll until it is fixed. +// +// Decoding over baseline rather than over conf is what makes the file a declarative overlay: keys +// it omits fall back to the startup values instead of inheriting the previous file's. +func applyConfigFile(conf, baseline *Config, path string, last []byte) (read []byte, ok bool) { + data, err := os.ReadFile(path) + if err != nil { + // An optional ConfigMap that is not mounted yet is the normal steady state, not something + // to log once per interval forever. A file that goes away is deliberately not treated as a + // reset either: a transient read error would otherwise flap levels. Emptying the file to + // `{}` is the reset. + return nil, false + } + if bytes.Equal(data, last) { + return data, false + } + + var file liveLevels + decoder := yaml.NewDecoder(bytes.NewReader(data)) + decoder.KnownFields(true) + if err := decoder.Decode(&file); err != nil && !errors.Is(err, io.EOF) { + Warnw("could not parse log config (only level and component_levels can change at runtime), keeping the one in force", + err, "path", path) + return data, false + } + + next := baseline.snapshot() + if file.Level != nil { + next.Level = *file.Level + } + maps.Copy(next.ComponentLevels, file.ComponentLevels) + if err := conf.updateLevels(next.Level, next.ComponentLevels); err != nil { + Warnw("could not apply log config, keeping the one in force", err, "path", path) + return data, false + } + return data, true +} + +// liveLevels is what the file may set: the part of the logging block a running logger can take. +// Everything else there is read once when a logger is built, so the file is decoded strictly and +// one that names any other key is rejected rather than half-applied. +type liveLevels struct { + Level *string `yaml:"level"` + ComponentLevels map[string]string `yaml:"component_levels"` +} diff --git a/logger/configwatch_test.go b/logger/configwatch_test.go new file mode 100644 index 000000000..9a35cc9b1 --- /dev/null +++ b/logger/configwatch_test.go @@ -0,0 +1,219 @@ +package logger + +import ( + "os" + "path/filepath" + "testing" + "time" + + "github.com/stretchr/testify/require" + "go.uber.org/zap/zapcore" +) + +func writeFile(t *testing.T, path, body string) { + t.Helper() + require.NoError(t, os.WriteFile(path, []byte(body), 0o600)) +} + +func TestApplyConfigFile(t *testing.T) { + t.Run("a level change reaches a live logger", func(t *testing.T) { + conf := &Config{Level: "info"} + l, err := NewZapLogger(conf) + require.NoError(t, err) + core := zapLoggerCore(l) + require.False(t, core.Enabled(zapcore.DebugLevel)) + + path := filepath.Join(t.TempDir(), "logging.yaml") + writeFile(t, path, "level: debug\n") + + applied, ok := applyConfigFile(conf, conf.snapshot(), path, nil) + require.True(t, ok) + require.NotEmpty(t, applied) + require.True(t, zapLoggerCore(l).Enabled(zapcore.DebugLevel), + "the atomic level behind the existing logger must move, not just Config.Level") + }) + + t.Run("keys absent from the file keep their current values", func(t *testing.T) { + // Update assigns every field, so applying a partial file over a zero Config would wipe + // these. component_levels is where livekit-server lands pion_level. + conf := &Config{ + Level: "info", + JSON: true, + Sample: true, + SampleInitial: 7, + ComponentLevels: map[string]string{"pion": "error"}, + } + _, err := NewZapLogger(conf) + require.NoError(t, err) + + path := filepath.Join(t.TempDir(), "logging.yaml") + writeFile(t, path, "level: warn\n") + + _, ok := applyConfigFile(conf, conf.snapshot(), path, nil) + require.True(t, ok) + require.Equal(t, "warn", conf.Level) + require.True(t, conf.JSON) + require.True(t, conf.Sample) + require.Equal(t, 7, conf.SampleInitial) + require.Equal(t, map[string]string{"pion": "error"}, conf.ComponentLevels) + }) + + t.Run("a component level in the file merges with the existing ones", func(t *testing.T) { + conf := &Config{Level: "info", ComponentLevels: map[string]string{"pion": "error"}} + l, err := NewZapLogger(conf) + require.NoError(t, err) + + path := filepath.Join(t.TempDir(), "logging.yaml") + writeFile(t, path, "component_levels:\n psrpc: debug\n") + + _, ok := applyConfigFile(conf, conf.snapshot(), path, nil) + require.True(t, ok) + require.Equal(t, "error", conf.ComponentLevels["pion"]) + require.Equal(t, "debug", conf.ComponentLevels["psrpc"]) + require.True(t, zapLoggerCore(l.WithComponent("psrpc")).Enabled(zapcore.DebugLevel)) + }) + + t.Run("unchanged bytes are not reapplied", func(t *testing.T) { + conf := &Config{Level: "info"} + path := filepath.Join(t.TempDir(), "logging.yaml") + writeFile(t, path, "level: debug\n") + + applied, ok := applyConfigFile(conf, conf.snapshot(), path, nil) + require.True(t, ok) + _, ok = applyConfigFile(conf, conf.snapshot(), path, applied) + require.False(t, ok) + }) + + t.Run("malformed yaml keeps the last good config", func(t *testing.T) { + conf := &Config{Level: "info"} + l, err := NewZapLogger(conf) + require.NoError(t, err) + + path := filepath.Join(t.TempDir(), "logging.yaml") + writeFile(t, path, "level: [not, a, string\n") + + _, ok := applyConfigFile(conf, conf.snapshot(), path, nil) + require.False(t, ok) + require.Equal(t, "info", conf.Level) + require.False(t, zapLoggerCore(l).Enabled(zapcore.DebugLevel)) + }) + + t.Run("an unchanged malformed file is reported once, not on every poll", func(t *testing.T) { + tee, logs := observerTee() + l, err := NewZapLogger(&Config{Level: "info"}, tee) + require.NoError(t, err) + SetLogger(l, "TEST") + t.Cleanup(func() { + defaultLogger = LogRLogger(discardLogger) + pkgLogger = LogRLogger(discardLogger) + }) + + conf := &Config{Level: "info"} + path := filepath.Join(t.TempDir(), "logging.yaml") + writeFile(t, path, "level: [not, a, string\n") + + var last []byte + for range 3 { + last, _ = applyConfigFile(conf, conf.snapshot(), path, last) + } + require.Equal(t, 1, logs.FilterMessageSnippet("could not parse log config").Len()) + }) + + t.Run("a key a running logger cannot take rejects the file", func(t *testing.T) { + conf := &Config{Level: "info"} + path := filepath.Join(t.TempDir(), "logging.yaml") + writeFile(t, path, "level: debug\njson: true\n") + + _, ok := applyConfigFile(conf, conf.snapshot(), path, nil) + require.False(t, ok) + require.Equal(t, "info", conf.Level, "a rejected file must not be half-applied") + require.False(t, conf.JSON) + }) + + t.Run("a missing file is tolerated", func(t *testing.T) { + conf := &Config{Level: "info"} + _, ok := applyConfigFile(conf, conf.snapshot(), filepath.Join(t.TempDir(), "absent.yaml"), nil) + require.False(t, ok) + require.Equal(t, "info", conf.Level) + }) +} + +func TestWatchConfigFile(t *testing.T) { + conf := &Config{Level: "info"} + l, err := NewZapLogger(conf) + require.NoError(t, err) + + path := filepath.Join(t.TempDir(), "logging.yaml") + writeFile(t, path, "level: info\n") + + stop := WatchConfigFile(conf, path, 5*time.Millisecond) + t.Cleanup(stop) + + writeFile(t, path, "level: debug\n") + require.Eventually(t, func() bool { + return zapLoggerCore(l).Enabled(zapcore.DebugLevel) + }, 2*time.Second, 5*time.Millisecond, "watcher should pick up the rewritten file") + + stop() + writeFile(t, path, "level: error\n") + time.Sleep(50 * time.Millisecond) + require.True(t, zapLoggerCore(l).Enabled(zapcore.DebugLevel), "stop must end the polling") +} + +// Applying config while a running logger reads it: the root leveler resolves the levels through +// ResolveComponentLevel under Config's own lock, and WithItemSampler reads the item sampler +// settings with no lock at all, so the watcher may write the levels and nothing else. +// Meaningful under -race. +func TestApplyConfigFileWhileResolvingComponents(t *testing.T) { + conf := &Config{Level: "info", ComponentLevels: map[string]string{"pion": "error"}, ItemSampleSeconds: 1} + l, err := NewZapLogger(conf) + require.NoError(t, err) + + dir := t.TempDir() + path := filepath.Join(dir, "logging.yaml") + + done := make(chan struct{}) + go func() { + defer close(done) + for i := 0; i < 500; i++ { + _ = zapLoggerCore(l.WithComponent("psrpc").WithComponent("Egress")) + _ = l.WithItemSampler() + } + }() + + var last []byte + for i, level := range []string{"debug", "warn", "info", "error"} { + writeFile(t, path, "level: "+level+"\n") + applied, ok := applyConfigFile(conf, conf.snapshot(), path, last) + require.True(t, ok, "iteration %d", i) + last = applied + } + <-done + require.Equal(t, "error", conf.Level) + require.Equal(t, "error", conf.ComponentLevels["pion"]) +} + +// The reset path the chart documents: emptying the file must put the startup levels back, not +// leave the last override in force. Applying each file over a baseline rather than over the +// config currently in force is what makes this hold. +func TestEmptyFileRestoresStartupConfig(t *testing.T) { + conf := &Config{Level: "info", ComponentLevels: map[string]string{"pion": "error"}} + l, err := NewZapLogger(conf) + require.NoError(t, err) + baseline := conf.snapshot() + + path := filepath.Join(t.TempDir(), "logging.yaml") + writeFile(t, path, "level: debug\ncomponent_levels:\n psrpc: debug\n") + applied, ok := applyConfigFile(conf, baseline, path, nil) + require.True(t, ok) + require.Equal(t, "debug", conf.Level) + require.True(t, zapLoggerCore(l).Enabled(zapcore.DebugLevel)) + + writeFile(t, path, "{}\n") + _, ok = applyConfigFile(conf, baseline, path, applied) + require.True(t, ok) + require.Equal(t, "info", conf.Level) + require.Equal(t, map[string]string{"pion": "error"}, conf.ComponentLevels, + "a component the file no longer names must stop applying") + require.False(t, zapLoggerCore(l).Enabled(zapcore.DebugLevel)) +} diff --git a/logger/logger.go b/logger/logger.go index 0d1af7981..8d4af8211 100644 --- a/logger/logger.go +++ b/logger/logger.go @@ -198,6 +198,7 @@ func FromZapLogger(log *zap.Logger, conf *Config, opts ...ZapLoggerOption) (ZapL zc.root.Refresh() return nil }) + startConfigWatchFromEnv(conf) for _, opt := range opts { opt(zc) }