Files
releases/v1/plugins/logger_integration_test.go
Stephan Renatus bc8e23a174 logging: keep forwarding from BufferedLogger after Flush() (#8544)
The BufferedLogger introduced for logger plugins is created at startup
and passed to the `*plugins.Manager`. Plugins (bundle, discovery,
status, logs) cache `manager.Logger()` in a field at construction time.
After `Manager.Start()`, `ResolveBufferedLogger` flushes the buffer and
swaps the `Manager'`s logger to a `StandardLogger` — but the plugins
still hold the old `BufferedLogger`. Since bundle loading is async, the
"Bundle loaded and activated successfully" message (and similar) gets
written to the already-flushed buffer where nobody reads it.

Signed-off-by: Stephan Renatus <stephan.renatus@gmail.com>
2026-04-22 07:42:51 +00:00

186 lines
5.2 KiB
Go

package plugins_test
import (
"context"
"log/slog"
"testing"
"testing/synctest"
"time"
"github.com/open-policy-agent/opa/v1/logging"
"github.com/open-policy-agent/opa/v1/logging/test"
"github.com/open-policy-agent/opa/v1/plugins"
"github.com/open-policy-agent/opa/v1/storage/inmem"
)
type testLoggerPlugin struct {
manager *plugins.Manager
logger *test.Logger
}
const testLoggerName = "test_logger"
func (p *testLoggerPlugin) Start(context.Context) error {
p.manager.UpdatePluginStatus(testLoggerName, &plugins.Status{State: plugins.StateOK})
return nil
}
func (p *testLoggerPlugin) Stop(context.Context) {
p.manager.UpdatePluginStatus(testLoggerName, &plugins.Status{State: plugins.StateNotReady})
}
func (*testLoggerPlugin) Reconfigure(context.Context, any) {}
func (p *testLoggerPlugin) Logger() slog.Handler {
return logging.NewSlogHandler(p.logger)
}
type testLoggerFactory struct {
logger *test.Logger
}
func (*testLoggerFactory) Validate(manager *plugins.Manager, config []byte) (any, error) {
return nil, nil
}
func (f *testLoggerFactory) New(manager *plugins.Manager, config any) plugins.Plugin {
return &testLoggerPlugin{
manager: manager,
logger: f.logger,
}
}
func TestBufferedLoggerIntegration(t *testing.T) {
synctest.Test(t, func(t *testing.T) {
ctx := t.Context()
bufferedLogger := logging.NewBufferedLogger(1000)
startTime := time.Now()
bufferedLogger.Info("early log 1")
time.Sleep(10 * time.Millisecond)
bufferedLogger.Debug("early log 2")
bufferedLogger.WithFields(map[string]any{"key": "value"}).Warn("early log with fields")
bufferedLogger.Error("early error log")
testLog := test.New()
factory := &testLoggerFactory{logger: testLog}
config := []byte(`{"logger": {"plugin": "test_logger"}}`)
manager, err := plugins.New(config, "test-instance", inmem.New(),
plugins.Logger(bufferedLogger))
if err != nil {
t.Fatalf("Failed to create manager: %v", err)
}
manager.Register(testLoggerName, factory.New(manager, nil))
if err := manager.Init(ctx); err != nil {
t.Fatalf("Failed to initialize manager: %v", err)
}
if err := manager.Start(ctx); err != nil {
t.Fatalf("Failed to start manager: %v", err)
}
defer manager.Stop(ctx)
p := manager.Plugin(testLoggerName)
if p == nil {
t.Fatal("Logger plugin not found")
}
loggerPlugin, ok := p.(plugins.LoggerPlugin)
if !ok {
t.Fatal("Plugin does not implement LoggerPlugin interface")
}
handler := loggerPlugin.Logger()
// Wrap the slog.Handler in a Logger adapter
targetLogger := logging.NewLoggerFromSlogHandler(handler, logging.Debug)
bufferedLogger.Flush(targetLogger)
time.Sleep(100 * time.Millisecond)
entries := testLog.Entries()
if len(entries) < 3 {
t.Fatalf("Expected at least 3 buffered entries (Debug filtered out), got %d", len(entries))
}
expectedMessages := []string{"early log 1", "early log with fields", "early error log"}
for i, expected := range expectedMessages {
if entries[i].Message != expected {
t.Errorf("Entry %d: expected message %q, got %q", i, expected, entries[i].Message)
}
}
if entries[1].Fields == nil || entries[1].Fields["key"] != "value" {
t.Errorf("Entry 1: expected fields with key=value, got %v", entries[1].Fields)
}
for i, entry := range entries[:3] {
if entry.Time.Before(startTime) {
t.Errorf("Entry %d: log time %v is before start time %v", i, entry.Time, startTime)
}
}
beforeDirectLog := len(entries)
// After flush, use the target logger directly
targetLogger.Info("direct log after plugin set")
time.Sleep(50 * time.Millisecond)
entries = testLog.Entries()
if len(entries) != beforeDirectLog+1 {
t.Errorf("Expected %d entries after direct log, got %d", beforeDirectLog+1, len(entries))
}
lastEntry := entries[len(entries)-1]
if lastEntry.Message != "direct log after plugin set" {
t.Errorf("Last entry: expected 'direct log after plugin set', got %q", lastEntry.Message)
}
})
}
func TestBufferedLoggerWithFields(t *testing.T) {
bufferedLogger := logging.NewBufferedLogger(1000)
testLog := test.New()
logger1 := bufferedLogger.WithFields(map[string]any{"component": "runtime"})
logger1.Info("buffered message")
bufferedLogger.Flush(testLog)
// After flush, log via the SAME cached WithFields reference.
// This is the scenario where plugins hold a stale reference to the
// BufferedLogger: the forwarding behavior ensures these logs reach
// the resolved target.
logger1.Info("forwarded message")
// Also log directly on the target for comparison.
logger2 := testLog.WithFields(map[string]any{"component": "server"})
logger2.Info("direct message")
entries := testLog.Entries()
if len(entries) != 3 {
t.Fatalf("Expected 3 entries, got %d", len(entries))
}
if entries[0].Fields["component"] != "runtime" {
t.Errorf("First entry: expected component=runtime, got %v", entries[0].Fields)
}
if entries[1].Message != "forwarded message" {
t.Errorf("Second entry: expected 'forwarded message', got %q", entries[1].Message)
}
if entries[1].Fields["component"] != "runtime" {
t.Errorf("Second entry: expected component=runtime, got %v", entries[1].Fields)
}
if entries[2].Fields["component"] != "server" {
t.Errorf("Third entry: expected component=server, got %v", entries[2].Fields)
}
}