Skip to content

Commit a4f439e

Browse files
authored
feat: add FlightRecorder sink for deferred low-level logging (#237)
* feat: add Buffer sink for deferred low-level logging Buffer is a Sink that holds entries below a configured level in a fixed-size in-memory ring, forwarding entries at or above the level immediately. Flush replays the buffered entries to the wrapped sinks. This lets a program run its logger at LevelDebug to capture verbose entries while only emitting them on demand (e.g. on an error or connection failure), with a configurable buffer size. When full, the oldest entry is evicted. * feat: add Flusher interface for deferred sinks Flusher (Flush(ctx)) lets callers trigger a flush without depending on the concrete Buffer type, which is handy when the buffer is created in one place and flushed in another. Buffer implements it. * refactor: rename buffer.go and buffer_test.go to flightrecorder.go and flightrecorder_test.go * refactor: rename Buffer to FlightRecorder * refactor: make flight recording a Logger flag with auto-flush on error Rework the flight recorder from a standalone Buffer/Flusher sink into a Logger capability. Logger.FlightRecorder(size) enables a rolling in-memory history of below-level entries; the history is flushed to the sinks automatically when an entry at LevelError or above is logged, or manually via Logger.Flush. Enabling flight recording is independent of the Logger level, so it is order-independent with Leveled(). Derived loggers (With, Named, Leveled) share the underlying ring and capture their derived fields and names. Replace the fill-tracking branch in the ring with a monotonic written counter: the ring is allocated at full length and writes always land at written%size, so the write path is branchless and drain derives the entry count and oldest position from written. * refactor: track ring position and length instead of unbounded counter Replace the monotonic written counter with bounded pos and length fields, so neither grows without limit and there is no risk of an integer-overflow index. drain derives the oldest position from length and pos.
1 parent 244048f commit a4f439e

3 files changed

Lines changed: 402 additions & 0 deletions

File tree

‎flightrecorder.go‎

Lines changed: 114 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,114 @@
1+
package slog
2+
3+
import (
4+
"context"
5+
"sync"
6+
)
7+
8+
// FlightRecorder returns a Logger that records entries below its level in a
9+
// fixed-size, in-memory ring buffer instead of dropping them, while still
10+
// forwarding entries at or above its level to the sinks as usual. The recorded
11+
// entries are forwarded to the sinks ("flushed") automatically whenever an entry
12+
// at LevelError or above is logged, and can be flushed manually with
13+
// Logger.Flush.
14+
//
15+
// This keeps verbose (e.g. debug) logs available for diagnosing a failure
16+
// without emitting them during normal operation: run the Logger at LevelInfo
17+
// with a flight recorder, and the debug entries leading up to an error are
18+
// emitted only when the error occurs. When the buffer is full, the oldest
19+
// recorded entry is dropped to make room for the newest.
20+
//
21+
// Enabling flight recording is independent of the Logger's level, so
22+
// Make(sink).Leveled(LevelInfo).FlightRecorder(n) and
23+
// Make(sink).FlightRecorder(n).Leveled(LevelInfo) are equivalent. If size is
24+
// <= 0, no entries are recorded.
25+
//
26+
// The returned Logger and any Loggers derived from it share the same underlying
27+
// buffer.
28+
func (l Logger) FlightRecorder(size int) Logger {
29+
l.flightRecorder = newFlightRecorder(size)
30+
return l
31+
}
32+
33+
// Flush forwards the entries currently held by the Logger's flight recorder to
34+
// its sinks, oldest first, and empties the buffer. It is a no-op when flight
35+
// recording is not enabled. Flushing also happens automatically when an entry at
36+
// LevelError or above is logged.
37+
func (l Logger) Flush(ctx context.Context) {
38+
if l.flightRecorder != nil {
39+
l.flightRecorder.flush(ctx, l.sinks)
40+
}
41+
}
42+
43+
// flightRecorder keeps a rolling history of entries in a fixed-size ring buffer
44+
// until they are flushed. It is safe for concurrent use.
45+
type flightRecorder struct {
46+
capacity int
47+
48+
mu sync.Mutex
49+
ring []SinkEntry
50+
// pos is the index the next entry is written to; length is the number of
51+
// entries currently held. Both stay bounded by capacity, so neither grows
52+
// without limit.
53+
pos int
54+
length int
55+
}
56+
57+
func newFlightRecorder(size int) *flightRecorder {
58+
if size < 0 {
59+
size = 0
60+
}
61+
return &flightRecorder{
62+
capacity: size,
63+
ring: make([]SinkEntry, size),
64+
}
65+
}
66+
67+
// record stores e in the ring, overwriting the oldest entry once full.
68+
func (f *flightRecorder) record(e SinkEntry) {
69+
if f.capacity == 0 {
70+
return
71+
}
72+
f.mu.Lock()
73+
f.ring[f.pos] = e
74+
f.pos = (f.pos + 1) % f.capacity
75+
if f.length < f.capacity {
76+
f.length++
77+
}
78+
f.mu.Unlock()
79+
}
80+
81+
// flush forwards the recorded entries to sinks, oldest first, and empties the
82+
// buffer.
83+
func (f *flightRecorder) flush(ctx context.Context, sinks []Sink) {
84+
f.mu.Lock()
85+
entries := f.drain()
86+
f.mu.Unlock()
87+
88+
for _, e := range entries {
89+
for _, s := range sinks {
90+
s.LogEntry(ctx, e)
91+
}
92+
}
93+
}
94+
95+
// drain returns the recorded entries oldest first and resets the buffer. It must
96+
// be called with f.mu held.
97+
func (f *flightRecorder) drain() []SinkEntry {
98+
if f.length == 0 {
99+
return nil
100+
}
101+
// Once the ring is full, pos points at the oldest entry; before that the
102+
// oldest entry is at index 0.
103+
oldest := 0
104+
if f.length == f.capacity {
105+
oldest = f.pos
106+
}
107+
entries := make([]SinkEntry, 0, f.length)
108+
for i := 0; i < f.length; i++ {
109+
entries = append(entries, f.ring[(oldest+i)%f.capacity])
110+
}
111+
f.length = 0
112+
f.pos = 0
113+
return entries
114+
}

‎flightrecorder_test.go‎

Lines changed: 271 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,271 @@
1+
package slog_test
2+
3+
import (
4+
"context"
5+
"sync"
6+
"testing"
7+
8+
"cdr.dev/slog/v3"
9+
"cdr.dev/slog/v3/internal/assert"
10+
)
11+
12+
func TestFlightRecorder(t *testing.T) {
13+
t.Parallel()
14+
15+
t.Run("ForwardsAtOrAboveLevel", func(t *testing.T) {
16+
t.Parallel()
17+
18+
s := &fakeSink{}
19+
log := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(8)
20+
21+
log.Info(bg, "info")
22+
log.Warn(bg, "warn")
23+
24+
assert.Len(t, "forwarded", 2, s.entries)
25+
assert.Equal(t, "first", "info", s.entries[0].Message)
26+
assert.Equal(t, "second", "warn", s.entries[1].Message)
27+
})
28+
29+
t.Run("BuffersBelowLevelUntilFlush", func(t *testing.T) {
30+
t.Parallel()
31+
32+
s := &fakeSink{}
33+
log := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(8)
34+
35+
log.Debug(bg, "debug1")
36+
log.Debug(bg, "debug2")
37+
// Nothing below the level is forwarded yet.
38+
assert.Len(t, "held", 0, s.entries)
39+
40+
log.Flush(bg)
41+
assert.Len(t, "flushed", 2, s.entries)
42+
assert.Equal(t, "first", "debug1", s.entries[0].Message)
43+
assert.Equal(t, "second", "debug2", s.entries[1].Message)
44+
})
45+
46+
t.Run("InterleavesImmediateAndBuffered", func(t *testing.T) {
47+
t.Parallel()
48+
49+
s := &fakeSink{}
50+
log := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(8)
51+
52+
log.Debug(bg, "debug")
53+
log.Info(bg, "info")
54+
// The info entry is forwarded immediately; the debug entry waits.
55+
assert.Len(t, "immediate", 1, s.entries)
56+
assert.Equal(t, "immediate", "info", s.entries[0].Message)
57+
58+
log.Flush(bg)
59+
assert.Len(t, "after flush", 2, s.entries)
60+
assert.Equal(t, "buffered", "debug", s.entries[1].Message)
61+
})
62+
63+
t.Run("FlushesAutomaticallyBeforeError", func(t *testing.T) {
64+
t.Parallel()
65+
66+
s := &fakeSink{}
67+
log := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(8)
68+
69+
log.Debug(bg, "debug1")
70+
log.Debug(bg, "debug2")
71+
// The buffered history is still held until the error triggers a flush.
72+
assert.Len(t, "held", 0, s.entries)
73+
74+
log.Error(bg, "boom")
75+
// The buffered debug entries precede the error that triggered the flush.
76+
assert.Len(t, "flushed with error", 3, s.entries)
77+
assert.Equal(t, "first", "debug1", s.entries[0].Message)
78+
assert.Equal(t, "second", "debug2", s.entries[1].Message)
79+
assert.Equal(t, "error last", "boom", s.entries[2].Message)
80+
})
81+
82+
t.Run("CriticalFlushesBufferedHistory", func(t *testing.T) {
83+
t.Parallel()
84+
85+
s := &fakeSink{}
86+
log := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(8)
87+
88+
log.Debug(bg, "debug")
89+
log.Critical(bg, "critical")
90+
91+
assert.Len(t, "flushed", 2, s.entries)
92+
assert.Equal(t, "first", "debug", s.entries[0].Message)
93+
assert.Equal(t, "second", "critical", s.entries[1].Message)
94+
})
95+
96+
t.Run("EvictsOldestWhenFull", func(t *testing.T) {
97+
t.Parallel()
98+
99+
s := &fakeSink{}
100+
log := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(2)
101+
102+
log.Debug(bg, "debug1")
103+
log.Debug(bg, "debug2")
104+
log.Debug(bg, "debug3")
105+
106+
log.Flush(bg)
107+
// debug1 was evicted; the two newest remain in order.
108+
assert.Len(t, "flushed", 2, s.entries)
109+
assert.Equal(t, "first", "debug2", s.entries[0].Message)
110+
assert.Equal(t, "second", "debug3", s.entries[1].Message)
111+
})
112+
113+
t.Run("EvictsOldestAcrossMultipleWraps", func(t *testing.T) {
114+
t.Parallel()
115+
116+
s := &fakeSink{}
117+
log := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(2)
118+
119+
log.Debug(bg, "debug1")
120+
log.Debug(bg, "debug2")
121+
log.Debug(bg, "debug3")
122+
log.Debug(bg, "debug4")
123+
log.Debug(bg, "debug5")
124+
125+
log.Flush(bg)
126+
// Only the two newest survive after wrapping multiple times.
127+
assert.Len(t, "flushed", 2, s.entries)
128+
assert.Equal(t, "first", "debug4", s.entries[0].Message)
129+
assert.Equal(t, "second", "debug5", s.entries[1].Message)
130+
})
131+
132+
t.Run("FlushIsIdempotent", func(t *testing.T) {
133+
t.Parallel()
134+
135+
s := &fakeSink{}
136+
log := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(8)
137+
138+
log.Debug(bg, "debug")
139+
log.Flush(bg)
140+
log.Flush(bg)
141+
142+
assert.Len(t, "flushed once", 1, s.entries)
143+
})
144+
145+
t.Run("ZeroSizeDropsBelowLevel", func(t *testing.T) {
146+
t.Parallel()
147+
148+
s := &fakeSink{}
149+
log := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(0)
150+
151+
log.Debug(bg, "debug")
152+
log.Info(bg, "info")
153+
log.Flush(bg)
154+
155+
// Only the info entry is forwarded; the debug entry was dropped.
156+
assert.Len(t, "forwarded", 1, s.entries)
157+
assert.Equal(t, "only", "info", s.entries[0].Message)
158+
})
159+
160+
t.Run("OrderIndependentWithLeveled", func(t *testing.T) {
161+
t.Parallel()
162+
163+
s1 := &fakeSink{}
164+
s2 := &fakeSink{}
165+
leveledFirst := slog.Make(s1).Leveled(slog.LevelInfo).FlightRecorder(8)
166+
recorderFirst := slog.Make(s2).FlightRecorder(8).Leveled(slog.LevelInfo)
167+
168+
for _, log := range []slog.Logger{leveledFirst, recorderFirst} {
169+
log.Debug(bg, "debug")
170+
log.Info(bg, "info")
171+
}
172+
173+
// Both loggers hold the debug entry and forward the info entry.
174+
assert.Len(t, "leveled first immediate", 1, s1.entries)
175+
assert.Len(t, "recorder first immediate", 1, s2.entries)
176+
177+
leveledFirst.Flush(bg)
178+
recorderFirst.Flush(bg)
179+
180+
assert.Len(t, "leveled first flushed", 2, s1.entries)
181+
assert.Len(t, "recorder first flushed", 2, s2.entries)
182+
assert.Equal(t, "leveled first buffered", "debug", s1.entries[1].Message)
183+
assert.Equal(t, "recorder first buffered", "debug", s2.entries[1].Message)
184+
})
185+
186+
t.Run("DerivedLoggersShareRecorder", func(t *testing.T) {
187+
t.Parallel()
188+
189+
s := &fakeSink{}
190+
base := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(8)
191+
derived := base.Named("sub").With(slog.F("key", "value"))
192+
193+
derived.Debug(bg, "debug")
194+
// Flushing through the base logger emits the entry recorded via the
195+
// derived logger, and the derived fields and names are captured.
196+
base.Flush(bg)
197+
198+
assert.Len(t, "flushed", 1, s.entries)
199+
assert.Equal(t, "message", "debug", s.entries[0].Message)
200+
assert.Equal(t, "names", []string{"sub"}, s.entries[0].LoggerNames)
201+
assert.Equal(t, "fields", slog.M(slog.F("key", "value")), s.entries[0].Fields)
202+
})
203+
204+
t.Run("FlushForwardsToAllSinks", func(t *testing.T) {
205+
t.Parallel()
206+
207+
s1 := &fakeSink{}
208+
s2 := &fakeSink{}
209+
log := slog.Make(s1, s2).Leveled(slog.LevelInfo).FlightRecorder(8)
210+
211+
log.Debug(bg, "debug")
212+
log.Flush(bg)
213+
214+
assert.Len(t, "sink1", 1, s1.entries)
215+
assert.Len(t, "sink2", 1, s2.entries)
216+
})
217+
218+
t.Run("NoRecorderDropsBelowLevel", func(t *testing.T) {
219+
t.Parallel()
220+
221+
s := &fakeSink{}
222+
log := slog.Make(s).Leveled(slog.LevelInfo)
223+
224+
log.Debug(bg, "debug")
225+
log.Info(bg, "info")
226+
// Without a flight recorder, below-level entries are dropped and Flush
227+
// is a no-op.
228+
log.Flush(bg)
229+
230+
assert.Len(t, "forwarded", 1, s.entries)
231+
assert.Equal(t, "only", "info", s.entries[0].Message)
232+
})
233+
234+
t.Run("ConcurrentRecordAndFlush", func(t *testing.T) {
235+
t.Parallel()
236+
237+
s := &lockedSink{}
238+
log := slog.Make(s).Leveled(slog.LevelInfo).FlightRecorder(64)
239+
240+
var wg sync.WaitGroup
241+
for i := 0; i < 8; i++ {
242+
wg.Add(1)
243+
go func() {
244+
defer wg.Done()
245+
for j := 0; j < 100; j++ {
246+
log.Debug(bg, "debug")
247+
}
248+
log.Flush(bg)
249+
}()
250+
}
251+
wg.Wait()
252+
253+
// The exact count depends on interleaving; the test only asserts that
254+
// concurrent record and flush do not race or panic.
255+
log.Flush(bg)
256+
})
257+
}
258+
259+
// lockedSink is a concurrency-safe sink for race detection tests.
260+
type lockedSink struct {
261+
mu sync.Mutex
262+
entries []slog.SinkEntry
263+
}
264+
265+
func (s *lockedSink) LogEntry(_ context.Context, e slog.SinkEntry) {
266+
s.mu.Lock()
267+
defer s.mu.Unlock()
268+
s.entries = append(s.entries, e)
269+
}
270+
271+
func (s *lockedSink) Sync() {}

0 commit comments

Comments
 (0)