mirror of
https://github.com/open-policy-agent/opa.git
synced 2026-08-13 03:42:35 -06:00
bc8e23a174
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>
186 lines
5.2 KiB
Go
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)
|
|
}
|
|
}
|