-
-
Notifications
You must be signed in to change notification settings - Fork 4
Expand file tree
/
Copy pathmiddleware_threshold.go
More file actions
100 lines (85 loc) 路 2.81 KB
/
Copy pathmiddleware_threshold.go
File metadata and controls
100 lines (85 loc) 路 2.81 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
package slogsampling
import (
"context"
"time"
"log/slog"
slogmulti "github.com/samber/slog-multi"
"github.com/samber/slog-sampling/buffer"
)
type ThresholdSamplingOption struct {
// This will log the first `Threshold` log entries with the same hash,
// in a `Tick` interval as-is. Following that, it will allow `Rate` in the range [0.0, 1.0].
Tick time.Duration
Threshold uint64
Rate float64
// Group similar logs (default: by level and message)
Matcher Matcher
Buffer func(generator func(string) any) buffer.Buffer[string]
buffer buffer.Buffer[string]
// Optional hooks
OnAccepted func(context.Context, slog.Record)
OnDropped func(context.Context, slog.Record)
// When true, the first accepted record after a suppression window includes
// a "slog_sampling.dropped_count" attribute with the number of records that
// were dropped in the previous window. This gives operators visibility into
// suppression volume without a separate summary goroutine.
IncludeDroppedCount bool
}
// NewMiddleware returns a slog-multi middleware.
func (o ThresholdSamplingOption) NewMiddleware() slogmulti.Middleware {
if o.Rate < 0.0 || o.Rate > 1.0 {
panic("unexpected Rate: must be between 0.0 and 1.0")
}
if o.Matcher == nil {
o.Matcher = DefaultMatcher
}
if o.Buffer == nil {
o.Buffer = buffer.NewUnlimitedBuffer[string]()
}
o.buffer = o.Buffer(func(k string) any {
return newCounter()
})
return slogmulti.NewInlineMiddleware(
func(ctx context.Context, level slog.Level, next func(context.Context, slog.Level) bool) bool {
return next(ctx, level)
},
func(ctx context.Context, record slog.Record, next func(context.Context, slog.Record) error) error {
key := o.Matcher(ctx, &record)
v, _ := o.buffer.GetOrInsert(key)
cnt := v.(*counter)
n := cnt.Inc(o.Tick)
if n > o.Threshold {
// Fast path: skip expensive crypto/rand when Rate is 0 (drop all)
// or 1 (accept all). Only compute random when probabilistic sampling.
if o.Rate == 0 {
cnt.IncDropped()
hook(o.OnDropped, ctx, record)
return nil
}
if o.Rate < 1.0 {
random := randomPercentage()
if random >= o.Rate {
cnt.IncDropped()
hook(o.OnDropped, ctx, record)
return nil
}
}
}
// Attach dropped count from previous window to the first accepted
// record in a new window, so operators see suppression volume inline.
if o.IncludeDroppedCount && n == 1 {
if dropped := cnt.PrevDropped(); dropped > 0 {
record.AddAttrs(slog.Uint64("slog_sampling.dropped_count", dropped))
}
}
hook(o.OnAccepted, ctx, record)
return next(ctx, record)
},
func(attrs []slog.Attr, next func([]slog.Attr) slog.Handler) slog.Handler {
return next(attrs)
},
func(name string, next func(string) slog.Handler) slog.Handler {
return next(name)
},
)
}