diff --git a/collector/go.mod b/collector/go.mod index 216a5e1cc4a..2370c6eb2ff 100644 --- a/collector/go.mod +++ b/collector/go.mod @@ -352,6 +352,7 @@ require ( github.com/beorn7/perks v1.0.1 // indirect github.com/bitfield/gotestdox v0.2.2 // indirect github.com/blang/semver/v4 v4.0.0 // indirect + github.com/bluele/gcache v0.0.2 // indirect github.com/bmatcuk/doublestar/v4 v4.10.0 // indirect github.com/bodgit/plumbing v1.3.0 // indirect github.com/bodgit/sevenzip v1.6.1 // indirect @@ -386,6 +387,7 @@ require ( github.com/coreos/go-oidc/v3 v3.18.0 // indirect github.com/coreos/go-semver v0.3.1 // indirect github.com/coreos/go-systemd/v22 v22.7.0 // indirect + github.com/cornelk/hashmap v1.0.8 // indirect github.com/cyphar/filepath-securejoin v0.6.1 // indirect github.com/danieljoos/wincred v1.2.3 // indirect github.com/databricks/databricks-sql-go v1.11.0 // indirect @@ -845,6 +847,9 @@ require ( github.com/safchain/ethtool v0.7.0 // indirect github.com/sagikazarmark/locafero v0.11.0 // indirect github.com/samber/lo v1.53.0 // indirect + github.com/samber/slog-common v0.21.0 // indirect + github.com/samber/slog-multi v1.8.0 // indirect + github.com/samber/slog-sampling v1.6.0 // indirect github.com/samuel/go-zookeeper v0.0.0-20190923202752-2cc03de413da // indirect github.com/scaleway/scaleway-sdk-go v1.0.0-beta.36 // indirect github.com/sean-/seed v0.0.0-20170313163322-e2103e2c3529 // indirect diff --git a/collector/go.sum b/collector/go.sum index 8b112610a52..9c81c112578 100644 --- a/collector/go.sum +++ b/collector/go.sum @@ -633,6 +633,8 @@ github.com/bitfield/gotestdox v0.2.2 h1:x6RcPAbBbErKLnapz1QeAlf3ospg8efBsedU93CD github.com/bitfield/gotestdox v0.2.2/go.mod h1:D+gwtS0urjBrzguAkTM2wodsTQYFHdpx8eqRJ3N+9pY= github.com/blang/semver/v4 v4.0.0 h1:1PFHFE6yCCTv8C1TeyNNarDzntLi7wMI5i/pzqYIsAM= github.com/blang/semver/v4 v4.0.0/go.mod h1:IbckMUScFkM3pff0VJDNKRiT6TG/YpiHIM2yvyW5YoQ= +github.com/bluele/gcache v0.0.2 h1:WcbfdXICg7G/DGBh1PFfcirkWOQV+v077yF1pSy3DGw= +github.com/bluele/gcache v0.0.2/go.mod h1:m15KV+ECjptwSPxKhOhQoAFQVtUFjTVkc3H8o0t/fp0= github.com/bmatcuk/doublestar v1.1.1/go.mod h1:UD6OnuiIn0yFxxA2le/rnRU1G4RaI4UvFv1sNto9p6w= github.com/bmatcuk/doublestar/v4 v4.10.0 h1:zU9WiOla1YA122oLM6i4EXvGW62DvKZVxIe6TYWexEs= github.com/bmatcuk/doublestar/v4 v4.10.0/go.mod h1:xBQ8jztBU6kakFMg+8WGxn0c6z1fTSPVIjEY1Wr7jzc= @@ -737,6 +739,8 @@ github.com/coreos/go-systemd v0.0.0-20190321100706-95778dfbb74e/go.mod h1:F5haX7 github.com/coreos/go-systemd/v22 v22.7.0 h1:LAEzFkke61DFROc7zNLX/WA2i5J8gYqe0rSj9KI28KA= github.com/coreos/go-systemd/v22 v22.7.0/go.mod h1:xNUYtjHu2EDXbsxz1i41wouACIwT7Ybq9o0BQhMwD0w= github.com/coreos/pkg v0.0.0-20180928190104-399ea9e2e55f/go.mod h1:E3G3o1h8I7cfcXa63jLwjI0eiQQMgzzUDFVpN/nH/eA= +github.com/cornelk/hashmap v1.0.8 h1:nv0AWgw02n+iDcawr5It4CjQIAcdMMKRrs10HOJYlrc= +github.com/cornelk/hashmap v1.0.8/go.mod h1:RfZb7JO3RviW/rT6emczVuC/oxpdz4UsSB2LJSclR1k= github.com/cpuguy83/dockercfg v0.3.2 h1:DlJTyZGBDlXqUZ2Dk2Q3xHs/FtnooJJVaad2S9GKorA= github.com/cpuguy83/dockercfg v0.3.2/go.mod h1:sugsbF4//dDlL/i+S+rtpIWp+5h0BHJHfjj5/jFyUJc= github.com/cpuguy83/go-md2man v1.0.10/go.mod h1:SmD6nW6nTyfqj6ABTjUi3V3JVMnlJmwcJI5acqYI6dE= @@ -2237,6 +2241,12 @@ github.com/sagikazarmark/locafero v0.11.0 h1:1iurJgmM9G3PA/I+wWYIOw/5SyBtxapeHDc github.com/sagikazarmark/locafero v0.11.0/go.mod h1:nVIGvgyzw595SUSUE6tvCp3YYTeHs15MvlmU87WwIik= github.com/samber/lo v1.53.0 h1:t975lj2py4kJPQ6haz1QMgtId2gtmfktACxIXArw3HM= github.com/samber/lo v1.53.0/go.mod h1:4+MXEGsJzbKGaUEQFKBq2xtfuznW9oz/WrgyzMzRoM0= +github.com/samber/slog-common v0.21.0 h1:Wo2hTly1Br5RjYqX/BTWJJeDnTE85oWk/7vqlpZuAUc= +github.com/samber/slog-common v0.21.0/go.mod h1:d/6OaSlzdkl9PFpfRLgn8FwY1OW6EFmPtBpsHX4MrU0= +github.com/samber/slog-multi v1.8.0 h1:E05c1wnQ+8M58oQDBABlJ4TEIJWssNgtckso3zlaLlI= +github.com/samber/slog-multi v1.8.0/go.mod h1:6+3j/ILxDvAcLD75YdQAm6iKWu6AmwlohLgQxL/2aiI= +github.com/samber/slog-sampling v1.6.0 h1:ODK16Wse1139eo25P+APfYKQqLYE79LMnnTBcUe+OCA= +github.com/samber/slog-sampling v1.6.0/go.mod h1:2vMB0an9YwqxtBdzmcBOpO5ZoL6LI74ghC3jC4sJuk4= github.com/samuel/go-zookeeper v0.0.0-20190923202752-2cc03de413da h1:p3Vo3i64TCLY7gIfzeQaUJ+kppEO5WQG3cL8iE8tGHU= github.com/samuel/go-zookeeper v0.0.0-20190923202752-2cc03de413da/go.mod h1:gi+0XIa01GRL2eRQVjQkKGqKF3SF9vZR/HnPullcV2E= github.com/santhosh-tekuri/jsonschema/v6 v6.0.2 h1:KRzFb2m7YtdldCEkzs6KqmJw4nqEVZGK7IN2kJkjTuQ= diff --git a/docs/sources/reference/config-blocks/logging.md b/docs/sources/reference/config-blocks/logging.md index 0a29453f154..64474e874a2 100644 --- a/docs/sources/reference/config-blocks/logging.md +++ b/docs/sources/reference/config-blocks/logging.md @@ -65,6 +65,41 @@ Otherwise, `destination` defaults to `"stderr"`. {{< param "PRODUCT_NAME" >}} fails to start if `destination` is set to `"windows_event_log"` and {{< param "PRODUCT_NAME" >}} is not running on Windows. +## Blocks + +You can use the following blocks with `logging`: + +| Block | Description | Required | +| ----------------- | ---------------------------------------------- | -------- | +| [`rate_limiting`][rate_limiting] | Configure per-message rate limiting and sampling. | no | + +### `rate_limiting` + +The `rate_limiting` block enables per-message rate limiting and sampling of repeated log lines. + +| Name | Type | Description | Default | Required | +|------|------|-------------|---------|----------| +| `enabled` | `bool` | Enable per-message rate limiting. | `true` | no | +| `max_signatures` | `number` | Distinct signatures tracked; least recently used is evicted when full. | `1000` | no | +| `rate` | `number` | Fraction (0-1) of the over-threshold tail still admitted; `0` drops all excess. | `0` | no | +| `threshold` | `number` | Identical lines admitted per (component, level, message) per tick before sampling. | `10` | no | +| `tick` | `duration` | Sampling window. | `"10s"` | no | + +Rate limiting keys on the component, the log level, and the log message text (not attributes/fields). +Only identical repeated lines from the same component at the same level are throttled; distinct components/messages are independent (LRU-bounded by `max_signatures`). + +Log lines that share a constant message but differ only in attributes are treated as the same signature and throttled together. + +Log lines with an empty message, such as some `go-kit`-style logs emitted without a `msg` or `message` field, bypass rate limiting entirely and are always written. + +After suppression begins, the first admitted line of each new window carries a `slog_sampling.dropped_count` attribute. + +Dropped lines are counted by the `alloy_logging_suppressed_lines_total` metric (labeled by `level` and `component_id`). + +Set `enabled = false` to disable. Omitting the `rate_limiting` block leaves limiting enabled with defaults. + +[rate_limiting]: #rate_limiting + ## Retrieve logs You can retrieve the logs in different ways depending on your platform and installation method: diff --git a/go.mod b/go.mod index e26ceba1372..c3de486de40 100644 --- a/go.mod +++ b/go.mod @@ -969,6 +969,7 @@ require ( github.com/open-telemetry/opentelemetry-collector-contrib/pkg/translator/splunk v0.153.0 github.com/open-telemetry/opentelemetry-collector-contrib/processor/redactionprocessor v0.153.0 github.com/open-telemetry/opentelemetry-collector-contrib/receiver/nginxreceiver v0.153.0 + github.com/samber/slog-sampling v1.6.0 github.com/spf13/viper v1.21.0 github.com/vektah/gqlparser/v2 v2.5.36 github.com/zricethezav/gitleaks/v8 v8.30.1 @@ -1027,6 +1028,7 @@ require ( github.com/aws/aws-sdk-go-v2/service/rds v1.118.2 // indirect github.com/aws/aws-sdk-go-v2/service/signin v1.2.2 // indirect github.com/aymanbagabas/go-osc52/v2 v2.0.1 // indirect + github.com/bluele/gcache v0.0.2 // indirect github.com/bodgit/plumbing v1.3.0 // indirect github.com/bodgit/sevenzip v1.6.1 // indirect github.com/bodgit/windows v1.0.1 // indirect @@ -1045,6 +1047,7 @@ require ( github.com/coder/websocket v1.8.14 // indirect github.com/containerd/containerd/api v1.9.0 // indirect github.com/containerd/typeurl/v2 v2.2.3 // indirect + github.com/cornelk/hashmap v1.0.8 // indirect github.com/dsnet/compress v0.0.2-0.20230904184137-39efe44ab707 // indirect github.com/erikgeiser/coninput v0.0.0-20211004153227-1c3628e74d0f // indirect github.com/fatih/semgroup v1.2.0 // indirect @@ -1079,6 +1082,8 @@ require ( github.com/puzpuzpuz/xsync/v4 v4.5.0 // indirect github.com/rs/xid v1.6.0 // indirect github.com/sagikazarmark/locafero v0.11.0 // indirect + github.com/samber/slog-common v0.21.0 // indirect + github.com/samber/slog-multi v1.8.0 // indirect github.com/sijms/go-ora/v2 v2.9.0 // indirect github.com/sorairolake/lzip-go v0.3.8 // indirect github.com/sosodev/duration v1.4.0 // indirect diff --git a/go.sum b/go.sum index 814e33aef76..3a7aaede107 100644 --- a/go.sum +++ b/go.sum @@ -653,6 +653,8 @@ github.com/bitfield/gotestdox v0.2.2 h1:x6RcPAbBbErKLnapz1QeAlf3ospg8efBsedU93CD github.com/bitfield/gotestdox v0.2.2/go.mod h1:D+gwtS0urjBrzguAkTM2wodsTQYFHdpx8eqRJ3N+9pY= github.com/blang/semver/v4 v4.0.0 h1:1PFHFE6yCCTv8C1TeyNNarDzntLi7wMI5i/pzqYIsAM= github.com/blang/semver/v4 v4.0.0/go.mod h1:IbckMUScFkM3pff0VJDNKRiT6TG/YpiHIM2yvyW5YoQ= +github.com/bluele/gcache v0.0.2 h1:WcbfdXICg7G/DGBh1PFfcirkWOQV+v077yF1pSy3DGw= +github.com/bluele/gcache v0.0.2/go.mod h1:m15KV+ECjptwSPxKhOhQoAFQVtUFjTVkc3H8o0t/fp0= github.com/bmatcuk/doublestar v1.1.1/go.mod h1:UD6OnuiIn0yFxxA2le/rnRU1G4RaI4UvFv1sNto9p6w= github.com/bmatcuk/doublestar/v4 v4.10.0 h1:zU9WiOla1YA122oLM6i4EXvGW62DvKZVxIe6TYWexEs= github.com/bmatcuk/doublestar/v4 v4.10.0/go.mod h1:xBQ8jztBU6kakFMg+8WGxn0c6z1fTSPVIjEY1Wr7jzc= @@ -775,6 +777,8 @@ github.com/coreos/go-systemd/v22 v22.3.2/go.mod h1:Y58oyj3AT4RCenI/lSvhwexgC+NSV github.com/coreos/go-systemd/v22 v22.7.0 h1:LAEzFkke61DFROc7zNLX/WA2i5J8gYqe0rSj9KI28KA= github.com/coreos/go-systemd/v22 v22.7.0/go.mod h1:xNUYtjHu2EDXbsxz1i41wouACIwT7Ybq9o0BQhMwD0w= github.com/coreos/pkg v0.0.0-20180928190104-399ea9e2e55f/go.mod h1:E3G3o1h8I7cfcXa63jLwjI0eiQQMgzzUDFVpN/nH/eA= +github.com/cornelk/hashmap v1.0.8 h1:nv0AWgw02n+iDcawr5It4CjQIAcdMMKRrs10HOJYlrc= +github.com/cornelk/hashmap v1.0.8/go.mod h1:RfZb7JO3RviW/rT6emczVuC/oxpdz4UsSB2LJSclR1k= github.com/cpuguy83/dockercfg v0.3.2 h1:DlJTyZGBDlXqUZ2Dk2Q3xHs/FtnooJJVaad2S9GKorA= github.com/cpuguy83/dockercfg v0.3.2/go.mod h1:sugsbF4//dDlL/i+S+rtpIWp+5h0BHJHfjj5/jFyUJc= github.com/cpuguy83/go-md2man v1.0.10/go.mod h1:SmD6nW6nTyfqj6ABTjUi3V3JVMnlJmwcJI5acqYI6dE= @@ -2259,6 +2263,12 @@ github.com/sagikazarmark/locafero v0.11.0 h1:1iurJgmM9G3PA/I+wWYIOw/5SyBtxapeHDc github.com/sagikazarmark/locafero v0.11.0/go.mod h1:nVIGvgyzw595SUSUE6tvCp3YYTeHs15MvlmU87WwIik= github.com/samber/lo v1.53.0 h1:t975lj2py4kJPQ6haz1QMgtId2gtmfktACxIXArw3HM= github.com/samber/lo v1.53.0/go.mod h1:4+MXEGsJzbKGaUEQFKBq2xtfuznW9oz/WrgyzMzRoM0= +github.com/samber/slog-common v0.21.0 h1:Wo2hTly1Br5RjYqX/BTWJJeDnTE85oWk/7vqlpZuAUc= +github.com/samber/slog-common v0.21.0/go.mod h1:d/6OaSlzdkl9PFpfRLgn8FwY1OW6EFmPtBpsHX4MrU0= +github.com/samber/slog-multi v1.8.0 h1:E05c1wnQ+8M58oQDBABlJ4TEIJWssNgtckso3zlaLlI= +github.com/samber/slog-multi v1.8.0/go.mod h1:6+3j/ILxDvAcLD75YdQAm6iKWu6AmwlohLgQxL/2aiI= +github.com/samber/slog-sampling v1.6.0 h1:ODK16Wse1139eo25P+APfYKQqLYE79LMnnTBcUe+OCA= +github.com/samber/slog-sampling v1.6.0/go.mod h1:2vMB0an9YwqxtBdzmcBOpO5ZoL6LI74ghC3jC4sJuk4= github.com/samuel/go-zookeeper v0.0.0-20190923202752-2cc03de413da h1:p3Vo3i64TCLY7gIfzeQaUJ+kppEO5WQG3cL8iE8tGHU= github.com/samuel/go-zookeeper v0.0.0-20190923202752-2cc03de413da/go.mod h1:gi+0XIa01GRL2eRQVjQkKGqKF3SF9vZR/HnPullcV2E= github.com/satori/go.uuid v1.2.0/go.mod h1:dA0hQrYB0VpLJoorglMZABFdXlWrHn1NEOzdhQKdks0= diff --git a/internal/runtime/internal/controller/loader.go b/internal/runtime/internal/controller/loader.go index a118408832f..8c937e7f14a 100644 --- a/internal/runtime/internal/controller/loader.go +++ b/internal/runtime/internal/controller/loader.go @@ -29,8 +29,8 @@ import ( "github.com/grafana/alloy/internal/runtime/internal/worker" "github.com/grafana/alloy/internal/runtime/tracing" "github.com/grafana/alloy/internal/service" - "github.com/grafana/alloy/internal/util" astutil "github.com/grafana/alloy/internal/util/ast" + "github.com/grafana/alloy/internal/util/metricsutil" "github.com/grafana/alloy/syntax/ast" "github.com/grafana/alloy/syntax/diag" "github.com/grafana/alloy/syntax/vm" @@ -120,12 +120,12 @@ func NewLoader(opts LoaderOptions) (*Loader, error) { // These metrics already being registered indicates there's already a loader which exists for this controller. // Creating duplicate loaders should only happen in error states where we should not proceed further. One know // case of this is when remotecfg loads an invalid config and attempts to reload the prior config. - existing := util.MustRegisterOrReturnExisting(globals.Registerer, l.cc) + existing := metricsutil.MustRegisterOrReturnExisting(globals.Registerer, l.cc) if existing != nil { return nil, fmt.Errorf("a loader exists already exists for %q", globals.ControllerID) } - existing = util.MustRegisterOrReturnExisting(globals.Registerer, l.cm) + existing = metricsutil.MustRegisterOrReturnExisting(globals.Registerer, l.cm) if existing != nil { return nil, fmt.Errorf("a loader exists already exists for %q", globals.ControllerID) } diff --git a/internal/runtime/internal/controller/node_config_logging.go b/internal/runtime/internal/controller/node_config_logging.go index 8dc588cb45a..bfcd4779a03 100644 --- a/internal/runtime/internal/controller/node_config_logging.go +++ b/internal/runtime/internal/controller/node_config_logging.go @@ -25,6 +25,7 @@ type LoggingConfigNode struct { // NewLoggingConfigNode creates a new LoggingConfigNode from an initial ast.BlockStmt. // The underlying config isn't applied until Evaluate is called. func NewLoggingConfigNode(block *ast.BlockStmt, globals ComponentGlobals) *LoggingConfigNode { + globals.Logger.InitRateLimitMetrics(globals.Registerer) return &LoggingConfigNode{ nodeID: BlockComponentID(block).String(), componentName: block.GetBlockName(), @@ -38,6 +39,7 @@ func NewLoggingConfigNode(block *ast.BlockStmt, globals ComponentGlobals) *Loggi // NewDefaultLoggingConfigNode creates a new LoggingConfigNode with nil block and eval. // This will force evaluate to use the default logging options for this node. func NewDefaultLoggingConfigNode(globals ComponentGlobals) *LoggingConfigNode { + globals.Logger.InitRateLimitMetrics(globals.Registerer) return &LoggingConfigNode{ nodeID: loggingBlockID, componentName: loggingBlockID, diff --git a/internal/runtime/logging/deferred_handler.go b/internal/runtime/logging/deferred_handler.go index 1e3f29679b7..a1e9f7ec24e 100644 --- a/internal/runtime/logging/deferred_handler.go +++ b/internal/runtime/logging/deferred_handler.go @@ -88,9 +88,14 @@ func (d *deferredSlogHandler) buildHandlers(parent slog.Handler) { d.mut.Lock() defer d.mut.Unlock() - // Root node will not have attrs or groups. + // The root node has no attrs or groups. Route it through the Logger's + // persistent rootInjector, so the shared root handler (l.rlHolder), + // which may be rate-limited, sits between component loggers and the + // terminal handler. Reusing rootInjector, instead of building a new + // samplingInjector on every Update, avoids an allocation on config + // reloads that do not touch rate limiting. if parent == nil { - d.handle = d.l.handler + d.handle = d.l.rootInjector } else { if d.group != "" { d.handle = parent.WithGroup(d.group) diff --git a/internal/runtime/logging/logger.go b/internal/runtime/logging/logger.go index 80d9f6e503e..bc75fa725bb 100644 --- a/internal/runtime/logging/logger.go +++ b/internal/runtime/logging/logger.go @@ -8,6 +8,9 @@ import ( "sync" "time" + "go.uber.org/atomic" + + "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/common/model" "github.com/grafana/alloy/internal/component/common/loki" @@ -36,6 +39,42 @@ type Logger struct { // the optional Windows Event Log. handler *handler deferredSlog *deferredSlogHandler // Buffers slog output until config is loaded, then delegates to handler. + + // rootInjector is the samplingInjector at the root of the deferred + // handler tree. buildHandlers reuses this single instance on every + // Update, instead of building a fresh one, because the root injector + // needs only rlHolder and handler, both stable for the life of the + // Logger. Config changes reach it through rlHolder's version, not + // through rebuilding the injector itself. + rootInjector *samplingInjector + + // rlHolder holds the current root handler used by the samplingInjector + // in the deferred handler tree. This handler may be wrapped for + // rate-limit sampling. It starts as the bare terminal handler (rate + // limiting off), and Update swaps it atomically. The stored version + // increases on each swap, so samplingInjector instances know to rebuild + // their cached, per-component replay of the root handler. + rlHolder atomic.Pointer[versionedHandler] + // rlMetrics holds the suppressed-lines metric, or nil before + // InitRateLimitMetrics runs. buildRoot's OnDropped closure reads this + // pointer live on every drop, so it is an atomic.Pointer rather than a + // plain field: InitRateLimitMetrics can set it at any time, even after + // an earlier Update already built the root handler, and drops that + // follow will still be counted. + rlMetrics atomic.Pointer[rateLimitMetrics] + + // rlMut guards the rate-limiting block in Update (the rlApplied check, + // buildRoot call, version increase, and rlHolder store, which must + // happen as one unit) and InitRateLimitMetrics's one-time set of + // rlMetrics. Reading rlMetrics is lock-free (atomic); rlMut is not + // needed for that. + rlMut sync.Mutex + // rlApplied is the RateLimitingOptions last used to build the current + // rlHolder root. It is nil until the first Update runs. Update rebuilds + // the sampler, and bumps the stored version, only when rlApplied is nil + // or the new options differ from it. This way, a config reload that does + // not change rate limiting does not reset rate-limit budgets already in use. + rlApplied *RateLimitingOptions } var _ EnabledAware = (*Logger)(nil) @@ -94,6 +133,11 @@ func NewDeferred(w io.Writer) (*Logger, error) { writer: writer, handler: bh, } + // Rate limiting starts disabled: the injector's root is the bare + // terminal handler, so logging works as it did before rate limiting + // existed, until the first config Update enables it. + l.rlHolder.Store(&versionedHandler{version: 0, h: bh}) + l.rootInjector = newSamplingInjector(&l.rlHolder, l.handler) l.deferredSlog = newDeferredHandler(l) return l, nil @@ -115,6 +159,14 @@ func (l *Logger) Update(o Options) error { return fmt.Errorf("unrecognized log format %q", o.Format) } + rlOpts := defaultRateLimitingOptions() + if o.RateLimiting != nil { + rlOpts = *o.RateLimiting + } + if err := rlOpts.Validate(); err != nil { + return err + } + l.bufferMut.Lock() l.level.Set(slogLevel(o.Level).Level()) l.format.Set(o.Format) @@ -125,6 +177,16 @@ func (l *Logger) Update(o Options) error { l.writer.SetLokiWriter(o.WriteTo) l.bufferMut.Unlock() + l.rlMut.Lock() + if l.rlApplied == nil || *l.rlApplied != rlOpts { + root := buildRoot(rlOpts, l.handler, &l.rlMetrics) + next := l.rlHolder.Load().version + 1 + l.rlHolder.Store(&versionedHandler{version: next, h: root}) + applied := rlOpts + l.rlApplied = &applied + } + l.rlMut.Unlock() + // Rebuild deferred slog handlers outside bufferMut to avoid a deadlock // with concurrent Handle() calls (they hold a child handler's RLock // while waiting for bufferMut via addRecord). @@ -157,6 +219,18 @@ func (l *Logger) flushBuffer() { } } +// InitRateLimitMetrics sets up the suppressed-lines metric. Call this once. +// It takes effect right away, including for a logger that already has an +// Update call behind it: buildRoot's OnDropped closure reads rlMetrics live, +// so InitRateLimitMetrics does not need to run before the first Update. +func (l *Logger) InitRateLimitMetrics(reg prometheus.Registerer) { + l.rlMut.Lock() + defer l.rlMut.Unlock() + if l.rlMetrics.Load() == nil { + l.rlMetrics.Store(newRateLimitMetrics(reg)) + } +} + func (l *Logger) SetTemporaryWriter(w io.Writer) { l.writer.SetTemporaryWriter(w) } diff --git a/internal/runtime/logging/logger_event_log_test.go b/internal/runtime/logging/logger_event_log_test.go index 59356406445..75966b8258c 100644 --- a/internal/runtime/logging/logger_event_log_test.go +++ b/internal/runtime/logging/logger_event_log_test.go @@ -229,6 +229,11 @@ func TestUpdate_NoLossDuringConcurrentDestinationFlips(t *testing.T) { Level: LevelInfo, Format: FormatLogfmt, Destination: LogDestinationStderr, + // This test hammers an identical "hammer" message ~8000 times and + // asserts an exact delivery count; rate limiting (on by default) + // would suppress most of them under the same signature. Disable it + // so this remains a pure destination-flip delivery test. + RateLimiting: &RateLimitingOptions{Enabled: false}, })) sl := l.Slog() @@ -265,9 +270,10 @@ func TestUpdate_NoLossDuringConcurrentDestinationFlips(t *testing.T) { dest = LogDestinationWindowsEventLog } err := l.Update(Options{ - Level: LevelInfo, - Format: FormatLogfmt, - Destination: dest, + Level: LevelInfo, + Format: FormatLogfmt, + Destination: dest, + RateLimiting: &RateLimitingOptions{Enabled: false}, }) if err != nil { t.Errorf("Update failed: %v", err) @@ -294,9 +300,10 @@ func TestUpdate_NoLossDuringConcurrentDestinationFlips(t *testing.T) { // End the test in the default destination so the final accounting is // stable. require.NoError(t, l.Update(Options{ - Level: LevelInfo, - Format: FormatLogfmt, - Destination: LogDestinationStderr, + Level: LevelInfo, + Format: FormatLogfmt, + Destination: LogDestinationStderr, + RateLimiting: &RateLimitingOptions{Enabled: false}, })) innerLines := inner.Lines() diff --git a/internal/runtime/logging/logger_rl_test.go b/internal/runtime/logging/logger_rl_test.go new file mode 100644 index 00000000000..ac29648ce98 --- /dev/null +++ b/internal/runtime/logging/logger_rl_test.go @@ -0,0 +1,249 @@ +package logging + +import ( + "bytes" + "context" + "log/slog" + "strings" + "sync" + "testing" + "time" + + "github.com/prometheus/client_golang/prometheus" + "github.com/stretchr/testify/require" + "go.uber.org/goleak" +) + +func TestLoggerRateLimitEndToEnd(t *testing.T) { + defer goleak.VerifyNone(t) // PROOF: no goroutine leaked by this feature + var buf bytes.Buffer + // Tick is short, but long enough that the first burst of 10 calls lands + // in one window, even under -race. This lets the test also see the + // dropped-count annotation. slog-sampling adds this annotation only to + // the first admitted record of a new tick window, never to the window + // that produced the drops. + l, err := New(&buf, Options{ + Level: LevelInfo, Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: 100 * time.Millisecond, Threshold: 2, Rate: 0, MaxSignatures: 100}, + }) + require.NoError(t, err) + log := l.Slog() + for i := 0; i < 10; i++ { + log.Info("floody") + } + require.Equal(t, 2, strings.Count(buf.String(), "floody")) // Threshold admitted within the window + + time.Sleep(300 * time.Millisecond) // let the tick window roll over + log.Info("floody") + require.Equal(t, 3, strings.Count(buf.String(), "floody")) + require.Contains(t, buf.String(), "slog_sampling.dropped_count") // first admission of new window annotated with prior drops +} + +func TestLoggerDistinctComponentsNoCrossSuppress(t *testing.T) { + defer goleak.VerifyNone(t) + var buf bytes.Buffer + l, err := New(&buf, Options{Level: LevelInfo, Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: time.Hour, Threshold: 1, Rate: 0, MaxSignatures: 100}}) + require.NoError(t, err) + a := l.Slog().With("component_id", "a", "component_path", "/a") + b := l.Slog().With("component_id", "b", "component_path", "/b") + a.Info("same") + a.Info("same") // 2nd dropped + b.Info("same") // different component ⇒ own bucket ⇒ admitted + require.Equal(t, 2, strings.Count(buf.String(), "msg=same")) +} + +// TestLoggerSamePathDistinctComponentIDNoCrossSuppress checks a past bug, +// where compMatcher keyed only on component_path (the parent path, for +// example "/" for every top-level component), message, and level. Two +// top-level components with the same parent path but different +// component_id must not share a rate-limit bucket. +func TestLoggerSamePathDistinctComponentIDNoCrossSuppress(t *testing.T) { + defer goleak.VerifyNone(t) + var buf bytes.Buffer + l, err := New(&buf, Options{Level: LevelInfo, Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: time.Hour, Threshold: 1, Rate: 0, MaxSignatures: 100}}) + require.NoError(t, err) + a := l.Slog().With("component_path", "/", "component_id", "comp.a") + b := l.Slog().With("component_path", "/", "component_id", "comp.b") + a.Info("same") + a.Info("same") // 2nd dropped: same component, same signature + b.Info("same") // different component_id, same path ⇒ own bucket ⇒ admitted + require.Equal(t, 2, strings.Count(buf.String(), "msg=same")) +} + +// ctxVariantKey is an unexported context key used solely to build a non- +// context.Background() ctx for TestInjectorComponentKeyingViaContextVariant. +type ctxVariantKey struct{} + +// TestInjectorComponentKeyingViaContextVariant checks Handle's fallback +// path: when Handle gets a ctx other than context.Background(), for example +// from slog's *Context logging methods, it must still add component +// identity through withComponent(ctx, s.comp), not the cached bgCtx. This +// keeps distinct components with the same signature from suppressing each +// other. +func TestInjectorComponentKeyingViaContextVariant(t *testing.T) { + defer goleak.VerifyNone(t) + var buf bytes.Buffer + l, err := New(&buf, Options{Level: LevelInfo, Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: time.Hour, Threshold: 1, Rate: 0, MaxSignatures: 100}}) + require.NoError(t, err) + a := l.Slog().With("component_path", "/", "component_id", "comp.a") + b := l.Slog().With("component_path", "/", "component_id", "comp.b") + + // Deliberately not context.Background(), to exercise Handle's + // withComponent(ctx, s.comp) fallback rather than the cached bgCtx. + ctx := context.WithValue(context.Background(), ctxVariantKey{}, 1) + require.NotEqual(t, context.Background(), ctx) + + a.Log(ctx, slog.LevelInfo, "same") + a.Log(ctx, slog.LevelInfo, "same") // 2nd dropped: same component, same signature + b.Log(ctx, slog.LevelInfo, "same") // different component_id ⇒ own bucket ⇒ admitted despite the shared non-Background ctx + require.Equal(t, 2, strings.Count(buf.String(), "msg=same")) +} + +func TestLoggerDisabledByConfig(t *testing.T) { + defer goleak.VerifyNone(t) + var buf bytes.Buffer + l, err := New(&buf, Options{Level: LevelInfo, Format: FormatLogfmt, RateLimiting: &RateLimitingOptions{Enabled: false}}) + require.NoError(t, err) + log := l.Slog() + for i := 0; i < 10; i++ { + log.Info("noisy") + } + require.Equal(t, 10, strings.Count(buf.String(), "noisy")) +} + +func TestUpdateInvalidRateLimitingLeavesStateUnchanged(t *testing.T) { + defer goleak.VerifyNone(t) + var buf bytes.Buffer + l, err := New(&buf, Options{Level: LevelInfo, Format: FormatLogfmt, RateLimiting: &RateLimitingOptions{Enabled: false}}) + require.NoError(t, err) + + err = l.Update(Options{ + Level: LevelError, + Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: 0}, + }) + require.Error(t, err) + + // The level must not have been mutated: an Info record should still pass, + // proving the invalid RateLimiting config was rejected before any other + // state (level, format, writer) was applied. + require.True(t, l.Enabled(context.Background(), slog.LevelInfo)) + l.Slog().Info("still-info") + require.Contains(t, buf.String(), "still-info") +} + +// TestUpdateSameOptionsPreservesBudget checks that re-applying the same +// RateLimitingOptions on every Update, for example from an unrelated config +// reload, does not reset the sampler's per-signature counters. Before this +// fix, Update always rebuilt the sampler with a fresh LRU and counters on +// every call, so the repeated call below would wrongly re-admit the +// already-throttled line. +func TestUpdateSameOptionsPreservesBudget(t *testing.T) { + defer goleak.VerifyNone(t) + var buf bytes.Buffer + opts := Options{Level: LevelInfo, Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: time.Hour, Threshold: 2, Rate: 0, MaxSignatures: 100}} + l, err := New(&buf, opts) + require.NoError(t, err) + log := l.Slog() + + log.Info("steady") + log.Info("steady") + log.Info("steady") // 3rd: over threshold, dropped + require.Equal(t, 2, strings.Count(buf.String(), "msg=steady")) + + // Re-apply the identical options (simulating an unrelated config + // reload triggering LoggingConfigNode.Evaluate -> Update again). + require.NoError(t, l.Update(opts)) + + log.Info("steady") // still over the ORIGINAL window's budget: must stay dropped + require.Equal(t, 2, strings.Count(buf.String(), "msg=steady")) +} + +// TestUpdateChangedOptionsRebuilds confirms that when rate_limiting options +// change across an Update call, the new configuration takes effect for an +// existing logger. A change must still trigger a rebuild, not just skip it +// like unchanged options do. +func TestUpdateChangedOptionsRebuilds(t *testing.T) { + defer goleak.VerifyNone(t) + var buf bytes.Buffer + l, err := New(&buf, Options{Level: LevelInfo, Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: time.Hour, Threshold: 2, Rate: 0, MaxSignatures: 100}}) + require.NoError(t, err) + log := l.Slog() + + require.NoError(t, l.Update(Options{Level: LevelInfo, Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: time.Hour, Threshold: 1, Rate: 0, MaxSignatures: 100}})) + + log.Info("changed") + log.Info("changed") // 2nd: over the NEW threshold of 1, dropped + require.Equal(t, 1, strings.Count(buf.String(), "msg=changed")) +} + +// TestMetricInitAfterUpdateStillCounts checks that InitRateLimitMetrics +// takes effect even when it runs after the first Update already built the +// root handler. Before the fix, buildRoot closed over the *rateLimitMetrics +// value at build time; a nil value at that point (metrics not yet +// initialized) meant the counter stayed silently disabled forever, no +// matter when InitRateLimitMetrics ran afterward. +func TestMetricInitAfterUpdateStillCounts(t *testing.T) { + defer goleak.VerifyNone(t) + var buf bytes.Buffer + // New runs the first Update with no metrics registered yet. + l, err := New(&buf, Options{Level: LevelInfo, Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: time.Hour, Threshold: 1, Rate: 0, MaxSignatures: 100}}) + require.NoError(t, err) + + // Metrics are initialized only now, after the root handler already exists. + reg := prometheus.NewRegistry() + l.InitRateLimitMetrics(reg) + + log := l.Slog() + log.Info("late-metric") + log.Info("late-metric") // 2nd dropped: over threshold + + mfs, err := reg.Gather() + require.NoError(t, err) + var total float64 + for _, mf := range mfs { + if mf.GetName() != "alloy_logging_suppressed_lines_total" { + continue + } + for _, m := range mf.GetMetric() { + total += m.GetCounter().GetValue() + } + } + require.GreaterOrEqual(t, total, float64(1), "suppressed-lines counter must increment even when InitRateLimitMetrics runs after the first Update") +} + +// TestConcurrentUpdatesNoRace calls Update from many goroutines at once, +// with valid options that may differ. It checks that -race stays clean, +// nothing panics, and the logger still works afterward. +func TestConcurrentUpdatesNoRace(t *testing.T) { + defer goleak.VerifyNone(t) + var buf bytes.Buffer + l, err := New(&buf, Options{Level: LevelInfo, Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: time.Hour, Threshold: 10, Rate: 0, MaxSignatures: 100}}) + require.NoError(t, err) + + var wg sync.WaitGroup + for g := 0; g < 8; g++ { + wg.Add(1) + go func(g int) { + defer wg.Done() + for i := 0; i < 5; i++ { + threshold := uint64(1 + (g+i)%5) + err := l.Update(Options{Level: LevelInfo, Format: FormatLogfmt, + RateLimiting: &RateLimitingOptions{Enabled: true, Tick: time.Hour, Threshold: threshold, Rate: 0, MaxSignatures: 100}}) + require.NoError(t, err) + } + }(g) + } + wg.Wait() + + l.Slog().Info("post-concurrent-update") + require.Contains(t, buf.String(), "post-concurrent-update") +} diff --git a/internal/runtime/logging/options.go b/internal/runtime/logging/options.go index 1b2102734cc..71286056fa9 100644 --- a/internal/runtime/logging/options.go +++ b/internal/runtime/logging/options.go @@ -5,6 +5,7 @@ import ( "fmt" "log/slog" "math" + "time" "github.com/grafana/alloy/internal/component/common/loki" "github.com/grafana/alloy/syntax" @@ -16,7 +17,8 @@ type Options struct { Format Format `alloy:"format,attr,optional"` Destination LogDestination `alloy:"destination,attr,optional"` - WriteTo []loki.LogsReceiver `alloy:"write_to,attr,optional"` + WriteTo []loki.LogsReceiver `alloy:"write_to,attr,optional"` + RateLimiting *RateLimitingOptions `alloy:"rate_limiting,block,optional"` } // LogDestination is where to send the primary log output. @@ -49,13 +51,20 @@ func defaultDestination() LogDestination { return LogDestinationStderr } +// defaultRateLimitingOptions returns the default rate-limiting configuration. +func defaultRateLimitingOptions() RateLimitingOptions { + return RateLimitingOptions{Enabled: true, Tick: 10 * time.Second, Threshold: 10, Rate: 0, MaxSignatures: 1000} +} + // defaultOptions builds a fresh set of Logger defaults, evaluating the // platform-appropriate destination at call time. func defaultOptions() Options { + rl := defaultRateLimitingOptions() return Options{ - Level: LevelDefault, - Format: FormatDefault, - Destination: defaultDestination(), + Level: LevelDefault, + Format: FormatDefault, + Destination: defaultDestination(), + RateLimiting: &rl, } } @@ -158,3 +167,41 @@ func (ll *Format) UnmarshalText(text []byte) error { } return nil } + +// RateLimitingOptions configures log rate limiting per component and +// message. It is backed by github.com/samber/slog-sampling and is enabled +// by default. +type RateLimitingOptions struct { + Enabled bool `alloy:"enabled,attr,optional"` + Tick time.Duration `alloy:"tick,attr,optional"` + Threshold uint64 `alloy:"threshold,attr,optional"` + Rate float64 `alloy:"rate,attr,optional"` + MaxSignatures int `alloy:"max_signatures,attr,optional"` +} + +var _ syntax.Defaulter = (*RateLimitingOptions)(nil) + +// SetToDefault implements syntax.Defaulter. +func (o *RateLimitingOptions) SetToDefault() { + *o = defaultRateLimitingOptions() +} + +var _ syntax.Validator = (*RateLimitingOptions)(nil) + +// Validate implements syntax.Validator. +func (o RateLimitingOptions) Validate() error { + if !o.Enabled { + return nil + } + switch { + case o.Tick <= 0: + return fmt.Errorf("logging rate_limiting.tick must be > 0, got %v", o.Tick) + case o.Threshold == 0: + return fmt.Errorf("logging rate_limiting.threshold must be > 0") + case math.IsNaN(o.Rate) || o.Rate < 0 || o.Rate > 1: + return fmt.Errorf("logging rate_limiting.rate must be in [0,1], got %v", o.Rate) + case o.MaxSignatures <= 0: + return fmt.Errorf("logging rate_limiting.max_signatures must be > 0, got %d", o.MaxSignatures) + } + return nil +} diff --git a/internal/runtime/logging/options_test.go b/internal/runtime/logging/options_test.go index 3d1482a5dd6..b02cd5d65da 100644 --- a/internal/runtime/logging/options_test.go +++ b/internal/runtime/logging/options_test.go @@ -1,7 +1,9 @@ package logging import ( + "math" "testing" + "time" "github.com/grafana/alloy/syntax" "github.com/stretchr/testify/require" @@ -88,3 +90,32 @@ func TestOptions_EndToEnd(t *testing.T) { }) } } + +func TestOptionsDefaultEnablesRateLimiting(t *testing.T) { + var o Options + o.SetToDefault() + require.NotNil(t, o.RateLimiting) + require.True(t, o.RateLimiting.Enabled) + require.Equal(t, 10*time.Second, o.RateLimiting.Tick) + require.Equal(t, uint64(10), o.RateLimiting.Threshold) + require.Equal(t, 0.0, o.RateLimiting.Rate) + require.Equal(t, 1000, o.RateLimiting.MaxSignatures) +} + +func TestRateLimitingValidate(t *testing.T) { + valid := RateLimitingOptions{Enabled: true, Tick: time.Second, Threshold: 10, Rate: 0, MaxSignatures: 1000} + require.NoError(t, valid.Validate()) + require.NoError(t, RateLimitingOptions{Enabled: false}.Validate()) + for _, m := range []func(*RateLimitingOptions){ + func(o *RateLimitingOptions) { o.Tick = 0 }, + func(o *RateLimitingOptions) { o.Threshold = 0 }, + func(o *RateLimitingOptions) { o.Rate = -0.1 }, + func(o *RateLimitingOptions) { o.Rate = 1.1 }, + func(o *RateLimitingOptions) { o.Rate = math.NaN() }, + func(o *RateLimitingOptions) { o.MaxSignatures = 0 }, + } { + bad := valid + m(&bad) + require.Error(t, bad.Validate()) + } +} diff --git a/internal/runtime/logging/rl_bench_test.go b/internal/runtime/logging/rl_bench_test.go new file mode 100644 index 00000000000..36ad97f002a --- /dev/null +++ b/internal/runtime/logging/rl_bench_test.go @@ -0,0 +1,122 @@ +package logging + +import ( + "io" + "strconv" + "testing" + "time" + + "github.com/prometheus/client_golang/prometheus" +) + +// newRLBenchLogger builds a real *Logger that writes to io.Discard, using +// the given RateLimitingOptions. New builds the root handler first, and +// InitRateLimitMetrics runs after, the same order production code uses. The +// drop path exercises the live alloy_logging_suppressed_lines_total counter, +// the same as in production. +func newRLBenchLogger(b *testing.B, rl RateLimitingOptions) *Logger { + b.Helper() + l, err := New(io.Discard, Options{ + Level: LevelInfo, + Format: FormatLogfmt, + RateLimiting: &rl, + }) + if err != nil { + b.Fatalf("failed to create logger: %v", err) + } + l.InitRateLimitMetrics(prometheus.NewRegistry()) + return l +} + +// BenchmarkRL_Disabled measures the baseline passthrough path with rate +// limiting disabled entirely. +func BenchmarkRL_Disabled(b *testing.B) { + l := newRLBenchLogger(b, RateLimitingOptions{Enabled: false}) + log := l.Slog() + + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + log.Info("bench message") + } +} + +// BenchmarkRL_AdmitHot measures the steady-state admit path. The threshold +// is effectively unbounded, so every call is admitted: build the matcher +// key, increase the counter, and write to the terminal. +func BenchmarkRL_AdmitHot(b *testing.B) { + l := newRLBenchLogger(b, RateLimitingOptions{ + Enabled: true, + Tick: time.Hour, + Threshold: 1 << 62, + Rate: 0, + MaxSignatures: 1000, + }) + log := l.Slog() + + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + log.Info("bench message") + } +} + +// BenchmarkRL_DropHot measures the drop path. After the first call, every +// identical call is dropped: matcher, counter, and OnDropped metric, but no +// terminal write. +func BenchmarkRL_DropHot(b *testing.B) { + l := newRLBenchLogger(b, RateLimitingOptions{ + Enabled: true, + Tick: time.Hour, + Threshold: 1, + Rate: 0, + MaxSignatures: 1000, + }) + log := l.Slog() + + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + log.Info("bench message") + } +} + +// BenchmarkRL_Churn logs a different message each iteration so every call +// is a new signature, exercising the LRU insert path. +func BenchmarkRL_Churn(b *testing.B) { + l := newRLBenchLogger(b, RateLimitingOptions{ + Enabled: true, + Tick: time.Hour, + Threshold: 1 << 62, + Rate: 0, + MaxSignatures: 1000, + }) + log := l.Slog() + + b.ReportAllocs() + b.ResetTimer() + for i := 0; i < b.N; i++ { + log.Info("msg " + strconv.Itoa(i&0xffff)) + } +} + +// BenchmarkRL_Parallel measures lock/contention on the shared buffer when N +// goroutines log the same admit-hot line concurrently. +func BenchmarkRL_Parallel(b *testing.B) { + l := newRLBenchLogger(b, RateLimitingOptions{ + Enabled: true, + Tick: time.Hour, + Threshold: 1 << 62, + Rate: 0, + MaxSignatures: 1000, + }) + log := l.Slog() + + b.ReportAllocs() + b.ResetTimer() + b.RunParallel(func(pb *testing.PB) { + for pb.Next() { + log.Info("bench message") + } + }) +} diff --git a/internal/runtime/logging/sampling.go b/internal/runtime/logging/sampling.go new file mode 100644 index 00000000000..c2a706cb66e --- /dev/null +++ b/internal/runtime/logging/sampling.go @@ -0,0 +1,313 @@ +package logging + +import ( + "context" + "log/slog" + + "go.uber.org/atomic" + + "github.com/prometheus/client_golang/prometheus" + slogsampling "github.com/samber/slog-sampling" + "github.com/samber/slog-sampling/buffer" + + "github.com/grafana/alloy/internal/util/metricsutil" +) + +type componentInfo struct{ id, path string } + +// sniffComponent reads component_id, controller_id, component_path, and +// controller_path from attrs and merges them onto base. It returns the +// result. If component_id is set, it wins over controller_id; the same +// precedence applies to component_path over controller_path. +func sniffComponent(base componentInfo, attrs []slog.Attr) componentInfo { + c := base + for _, a := range attrs { + switch a.Key { + case "component_id": + c.id = a.Value.String() + case "controller_id": + if c.id == "" { + c.id = a.Value.String() + } + case "component_path": + c.path = a.Value.String() + case "controller_path": + if c.path == "" { + c.path = a.Value.String() + } + } + } + return c +} + +type ctxKey struct{} + +// withComponent stores c in ctx. A Matcher can read it back later. This +// avoids passing c through every log call site. +func withComponent(ctx context.Context, c componentInfo) context.Context { + return context.WithValue(ctx, ctxKey{}, c) +} + +// componentFromCtx returns the componentInfo stored by withComponent. It +// returns the zero value if none was stored. +func componentFromCtx(ctx context.Context) componentInfo { + c, _ := ctx.Value(ctxKey{}).(componentInfo) + return c +} + +// compMatcher is the slog-sampling Matcher used to build the rate limiter's +// signature key. Two records share one signature, and so share one +// rate-limit budget, when they have the same component path, component ID, +// level, and message. +// +// component_path alone is not enough. It is the parent path (for example, +// "/" for every top-level component), so many components share it. Without +// component_id, different top-level components that log the same message at +// the same level would share one signature and suppress each other. +func compMatcher(ctx context.Context, r *slog.Record) string { + c := componentFromCtx(ctx) + return c.path + "\x00" + c.id + "\x00" + r.Level.String() + "\x00" + r.Message +} + +// levelString converts a slog.Level to the lowercase label used in the +// suppressed-lines metric. +func levelString(l slog.Level) string { + switch { + case l < slog.LevelInfo: + return "debug" + case l < slog.LevelWarn: + return "info" + case l < slog.LevelError: + return "warn" + default: + return "error" + } +} + +// rateLimitMetrics counts log lines dropped by the rate limiter, by level +// and component. A nil *rateLimitMetrics is valid: its methods do nothing. +// Callers do not need to check for a missing registerer. +type rateLimitMetrics struct { + suppressed *prometheus.CounterVec +} + +// newRateLimitMetrics registers the alloy_logging_suppressed_lines_total +// counter vector on reg. It returns nil if reg is nil, so callers that do +// not want metrics, such as tests, can safely pass a nil registerer. +func newRateLimitMetrics(reg prometheus.Registerer) *rateLimitMetrics { + if reg == nil { + return nil + } + cv := prometheus.NewCounterVec(prometheus.CounterOpts{ + Name: "alloy_logging_suppressed_lines_total", + Help: "Total log lines dropped by the logger's rate limiter, by level and component.", + }, []string{"level", "component_id"}) + if existing := metricsutil.MustRegisterOrReturnExisting(reg, cv); existing != nil { + cvExisting, ok := existing.(*prometheus.CounterVec) + if !ok { + panic("alloy_logging_suppressed_lines_total already registered with unexpected collector type") + } + cv = cvExisting + } + return &rateLimitMetrics{suppressed: cv} +} + +// onDropped is the OnDropped hook for slog-sampling. It increments the +// suppressed-lines counter for the record's level and component. +func (m *rateLimitMetrics) onDropped(ctx context.Context, r slog.Record) { + if m == nil { + return + } + m.suppressed.WithLabelValues(levelString(r.Level), componentFromCtx(ctx).id).Inc() +} + +// buildRoot wraps terminal with sampling when rate limiting is enabled. +// Otherwise it returns terminal unchanged. +// +// metrics is a live pointer, not a captured value: OnDropped reads +// metrics.Load() on every drop, instead of closing over one +// *rateLimitMetrics snapshot at build time. This lets InitRateLimitMetrics +// take effect even when it runs after the root handler was already built by +// an earlier Update call. +func buildRoot(o RateLimitingOptions, terminal slog.Handler, metrics *atomic.Pointer[rateLimitMetrics]) slog.Handler { + if !o.Enabled { + return terminal + } + opt := slogsampling.ThresholdSamplingOption{ + Tick: o.Tick, + Threshold: o.Threshold, + Rate: o.Rate, + Matcher: compMatcher, + Buffer: buffer.NewLRUBuffer[string](o.MaxSignatures), + OnDropped: func(ctx context.Context, r slog.Record) { + if metrics != nil { + if m := metrics.Load(); m != nil { + m.onDropped(ctx, r) + } + } + }, + IncludeDroppedCount: true, + } + return opt.NewMiddleware()(terminal) +} + +// replayOp records one WithAttrs or WithGroup call, so it can be replayed +// later on a new terminal handler. Only one of attrs or group is set. +type replayOp struct { + attrs []slog.Attr + group string +} + +// versionedHandler pairs a handler with a version number. Logger's rlHolder +// uses it to hold the current root handler; Update increases the version +// each time the rate-limiting config changes. samplingInjector's cache uses +// the same type to hold its replay of that root handler, keyed by the same +// version, so it can compare versions to know when the replay is stale and +// must be rebuilt. +type versionedHandler struct { + version uint64 + h slog.Handler +} + +// handlerBox wraps a slog.Handler so bareCache can store it in an +// atomic.Pointer. atomic.Pointer needs a concrete type, and slog.Handler is +// an interface, so the box supplies that concrete type. +type handlerBox struct { + h slog.Handler +} + +// samplingInjector is a slog.Handler between component loggers and the +// shared root handler, which may be rate-limited. It does three things: +// +// - It tracks component identity (comp) read from WithAttrs calls. The +// rate limiter's Matcher reads comp from the context to key by +// component. +// - It records every WithAttrs/WithGroup call as a replayOp. When the root +// handler changes (for example, a rate-limiting config change) or is +// bypassed (an empty-message record), it replays these calls on the new +// or bare terminal handler. This keeps the rendered output the same. +// - It bypasses the rate limiter for empty-message records. Without this, +// all empty-message records would share one signature. +type samplingInjector struct { + comp componentInfo + ops []replayOp + + holder *atomic.Pointer[versionedHandler] + bare slog.Handler // bare terminal handler (no per-component attrs), used to bypass sampling for empty-message records + + // bgCtx is context.Background() with comp already added by + // withComponent. It is computed once per injector, not on every Handle + // call, because slog.Logger.Info/Warn/Error always pass + // context.Background(). This is by far the most common case on the + // admit path; see Handle. + bgCtx context.Context + + cache atomic.Pointer[versionedHandler] + bareCache atomic.Pointer[handlerBox] +} + +// newSamplingInjector creates a samplingInjector. holder points to the +// current root handler, which may be wrapped for sampling. bare is the +// terminal handler used to bypass sampling for empty-message records. +func newSamplingInjector(holder *atomic.Pointer[versionedHandler], bare slog.Handler) *samplingInjector { + s := &samplingInjector{holder: holder, bare: bare} + s.bgCtx = withComponent(context.Background(), s.comp) + return s +} + +// replay applies a recorded sequence of WithAttrs/WithGroup calls to h, in +// order. The result is the same handler as calling them directly on h. +func replay(h slog.Handler, ops []replayOp) slog.Handler { + for _, op := range ops { + if op.group != "" { + h = h.WithGroup(op.group) + } else { + h = h.WithAttrs(op.attrs) + } + } + return h +} + +// clone returns a new samplingInjector. It shares this injector's holder and +// bare terminal, but gets its own copy of ops and fresh, empty caches: a new +// ops slice means any old cached replay is stale. +func (s *samplingInjector) clone() *samplingInjector { + // The new injector shares holder and bare, and starts with empty caches + // because ops will differ. clone itself does not change comp (WithAttrs + // changes it on the returned injector below), so bgCtx also carries over + // unchanged from the parent. + ns := &samplingInjector{comp: s.comp, holder: s.holder, bare: s.bare, bgCtx: s.bgCtx} + ns.ops = make([]replayOp, len(s.ops), len(s.ops)+1) + copy(ns.ops, s.ops) + return ns +} + +// WithAttrs returns a new handler with attrs bound. It also reads attrs for +// component identity, so the rate limiter can key by component, and records +// the call as a replayOp, so the terminal handler still renders it. +func (s *samplingInjector) WithAttrs(attrs []slog.Attr) slog.Handler { + ns := s.clone() + ns.comp = sniffComponent(s.comp, attrs) + ns.bgCtx = withComponent(context.Background(), ns.comp) + ns.ops = append(ns.ops, replayOp{attrs: attrs}) + return ns +} + +// WithGroup returns a new handler with name added as an open group. It +// records the call as a replayOp, so the terminal handler still renders it. +func (s *samplingInjector) WithGroup(name string) slog.Handler { + if name == "" { + return s + } + ns := s.clone() + ns.ops = append(ns.ops, replayOp{group: name}) + return ns +} + +// Enabled calls the bare terminal handler, not the root handler, which may +// be wrapped for sampling. Sampling only drops records in Handle; it never +// changes whether a level is enabled. Calling the bare handler here skips +// the sampling wrapper's per-call cost on every slog.Info/Warn/Error call. +// s.bare uses the shared LevelVar that Update changes, and level gating does +// not depend on component, so this is correct for every derived injector, +// not just the root one. +func (s *samplingInjector) Enabled(ctx context.Context, l slog.Level) bool { + return s.bare.Enabled(ctx, l) +} + +// Handle sends empty-message records straight to the bare terminal handler. +// This skips the rate limiter, because blank-message records from different +// call sites would otherwise share one signature. All other records go +// through the current root handler, which may be rate-limited. Handle adds +// the component identity to ctx so the Matcher can read it. +func (s *samplingInjector) Handle(ctx context.Context, r slog.Record) error { + if r.Message == "" { + // Skip the sampler: unrelated no-message records must not share one signature. + bh := s.bareCache.Load() + if bh == nil { + bh = &handlerBox{h: replay(s.bare, s.ops)} + s.bareCache.Store(bh) + } + return bh.h.Handle(ctx, r) + } + vh := s.holder.Load() + c := s.cache.Load() + if c == nil || c.version != vh.version { + c = &versionedHandler{version: vh.version, h: replay(vh.h, s.ops)} + s.cache.Store(c) + } + // slog.Logger.Info/Warn/Error always pass context.Background(). Reuse + // the precomputed component ctx for this common case, instead of an + // allocation from context.WithValue on every admitted line. For any + // other ctx, for example from InfoContext, inject fresh instead: it + // may carry values or a Done channel that must not be lost. + var cctx context.Context + if ctx == context.Background() { + cctx = s.bgCtx + } else { + cctx = withComponent(ctx, s.comp) + } + return c.h.Handle(cctx, r) +} + +var _ slog.Handler = (*samplingInjector)(nil) diff --git a/internal/runtime/logging/sampling_test.go b/internal/runtime/logging/sampling_test.go new file mode 100644 index 00000000000..ba0698a1c38 --- /dev/null +++ b/internal/runtime/logging/sampling_test.go @@ -0,0 +1,194 @@ +package logging + +import ( + "bytes" + "context" + "io" + "log/slog" + "strings" + "testing" + "time" + + "go.uber.org/atomic" + + slogsampling "github.com/samber/slog-sampling" + "github.com/samber/slog-sampling/buffer" + "github.com/stretchr/testify/require" +) + +// TestSpikeThresholdAdmitsThenDrops confirms the Threshold middleware wraps a +// terminal handler, keys records with our Matcher, and admits `Threshold` +// records per tick, then drops the rest (rate=0). +func TestSpikeThresholdAdmitsThenDrops(t *testing.T) { + var buf bytes.Buffer + terminal := slog.NewTextHandler(&buf, nil) + + opt := slogsampling.ThresholdSamplingOption{ + Tick: time.Hour, // one window for the whole test + Threshold: 3, + Rate: 0, + Matcher: func(ctx context.Context, r *slog.Record) string { return r.Message }, + Buffer: buffer.NewLRUBuffer[string](100), + } + h := opt.NewMiddleware()(terminal) + logger := slog.New(h) + for i := 0; i < 10; i++ { + logger.Info("spam") + } + require.Equal(t, 3, strings.Count(buf.String(), "spam"), "expected exactly Threshold admitted") +} + +func TestSniffComponent(t *testing.T) { + t.Run("component_id wins over controller_id", func(t *testing.T) { + c := sniffComponent(componentInfo{}, []slog.Attr{ + slog.String("controller_id", "ctrl-1"), + slog.String("component_id", "comp-1"), + }) + require.Equal(t, "comp-1", c.id) + }) + + t.Run("controller_id used when component_id absent", func(t *testing.T) { + c := sniffComponent(componentInfo{}, []slog.Attr{ + slog.String("controller_id", "ctrl-1"), + }) + require.Equal(t, "ctrl-1", c.id) + }) + + t.Run("component_id already set on base is not overridden by controller_id", func(t *testing.T) { + c := sniffComponent(componentInfo{id: "existing"}, []slog.Attr{ + slog.String("controller_id", "ctrl-1"), + }) + require.Equal(t, "existing", c.id) + }) + + t.Run("component_path captured", func(t *testing.T) { + c := sniffComponent(componentInfo{}, []slog.Attr{ + slog.String("component_path", "/foo/bar"), + }) + require.Equal(t, "/foo/bar", c.path) + }) +} + +// TestSniffComponentControllerPath checks a past bug, where sniffComponent +// ignored controller_path. Controller log lines carry controller_id and +// controller_path, not component_id and component_path. Without this case, +// every controller log line got path="", so two nested controllers with the +// same leaf controller_id under different parents shared one rate-limit +// bucket and cross-suppressed each other. +func TestSniffComponentControllerPath(t *testing.T) { + t.Run("controller_path used when component_path absent", func(t *testing.T) { + c := sniffComponent(componentInfo{}, []slog.Attr{ + slog.String("controller_id", "controller_id"), + slog.String("controller_path", "controller_path"), + }) + require.Equal(t, componentInfo{id: "controller_id", path: "controller_path"}, c) + }) + + t.Run("component_path wins over controller_path", func(t *testing.T) { + c := sniffComponent(componentInfo{}, []slog.Attr{ + slog.String("controller_path", "/ctrl"), + slog.String("component_path", "/comp"), + }) + require.Equal(t, "/comp", c.path) + }) + + t.Run("component_path wins over controller_path regardless of attr order", func(t *testing.T) { + c := sniffComponent(componentInfo{}, []slog.Attr{ + slog.String("component_path", "/comp"), + slog.String("controller_path", "/ctrl"), + }) + require.Equal(t, "/comp", c.path) + }) +} + +func TestCompMatcherKeysOnPathIDLevelMessage(t *testing.T) { + mk := func(path, id string, level slog.Level, msg string) string { + ctx := withComponent(context.Background(), componentInfo{path: path, id: id}) + r := slog.NewRecord(time.Time{}, level, msg, 0) + return compMatcher(ctx, &r) + } + + base := mk("/a", "comp.a", slog.LevelInfo, "hello") + + require.Equal(t, base, mk("/a", "comp.a", slog.LevelInfo, "hello"), "identical path/id/level/message should produce the same key") + require.NotEqual(t, base, mk("/b", "comp.a", slog.LevelInfo, "hello"), "different path should produce a different key") + require.NotEqual(t, base, mk("/a", "comp.b", slog.LevelInfo, "hello"), "different component id (same path) should produce a different key") + require.NotEqual(t, base, mk("/a", "comp.a", slog.LevelWarn, "hello"), "different level should produce a different key") + require.NotEqual(t, base, mk("/a", "comp.a", slog.LevelInfo, "goodbye"), "different message should produce a different key") +} + +func TestBuildRootDisabledReturnsTerminal(t *testing.T) { + terminal := slog.NewTextHandler(&bytes.Buffer{}, nil) + got := buildRoot(RateLimitingOptions{Enabled: false}, terminal, nil) + require.Same(t, terminal, got) +} + +func newTestInjector(t *testing.T, root slog.Handler, bare slog.Handler) *samplingInjector { + t.Helper() + var holder atomic.Pointer[versionedHandler] + holder.Store(&versionedHandler{version: 1, h: root}) + return newSamplingInjector(&holder, bare) +} + +func TestInjectorRendersAttrsAndGroups(t *testing.T) { + var buf bytes.Buffer + term := slog.NewTextHandler(&buf, nil) + inj := newTestInjector(t, term, term) // no sampling, just render + h := inj.WithAttrs([]slog.Attr{slog.String("component_id", "x")}).WithGroup("g").WithAttrs([]slog.Attr{slog.String("k", "v")}).(*samplingInjector) + require.Equal(t, "x", h.comp.id) + rec := slog.NewRecord(time.Unix(0, 0), slog.LevelInfo, "hi", 0) + require.NoError(t, h.Handle(context.Background(), rec)) + out := buf.String() + require.Contains(t, out, "component_id=x") + require.Contains(t, out, "g.k=v") // group-nested attr rendered natively +} + +func TestInjectorEmptyMessageBypassesSampler(t *testing.T) { + var termBuf bytes.Buffer + term := slog.NewTextHandler(&termBuf, nil) + // A root that drops everything, to prove empty-message records go to + // `bare` (term), not root. AbsoluteSamplingOption panics when Max == 0, + // so we build the drop-all root with ThresholdSamplingOption{Threshold: + // 0, Rate: 0} instead, which admits nothing. + dropAll := slogsampling.ThresholdSamplingOption{ + Tick: time.Hour, + Threshold: 0, + Rate: 0, + Matcher: func(ctx context.Context, r *slog.Record) string { return "" }, + Buffer: buffer.NewLRUBuffer[string](10), + }.NewMiddleware()(slog.NewTextHandler(&bytes.Buffer{}, nil)) + inj := newTestInjector(t, dropAll, term) + rec := slog.NewRecord(time.Unix(0, 0), slog.LevelInfo, "", 0) + rec.AddAttrs(slog.String("k", "v")) + require.NoError(t, inj.Handle(context.Background(), rec)) + require.Contains(t, termBuf.String(), "k=v") // empty-msg reached bare terminal +} + +// TestInjectorEnabledMatchesTerminalLevel checks that Enabled calls the bare +// terminal handler, whose leveler reflects the current, live-updatable +// level, not the sampling-wrapped root. Root reports everything enabled +// (LevelDebug), while bare is gated at Info. If Enabled called root instead +// of bare, the LevelDebug assertion below would fail. +func TestInjectorEnabledMatchesTerminalLevel(t *testing.T) { + root := slog.NewTextHandler(io.Discard, &slog.HandlerOptions{Level: slog.LevelDebug}) + bare := slog.NewTextHandler(io.Discard, &slog.HandlerOptions{Level: slog.LevelInfo}) + inj := newTestInjector(t, root, bare) + + require.False(t, inj.Enabled(context.Background(), slog.LevelDebug), "Enabled must reflect the terminal's level, not the sampling root's") + require.True(t, inj.Enabled(context.Background(), slog.LevelInfo)) +} + +func TestInjectorReDerivesOnVersionBump(t *testing.T) { + var bufA, bufB bytes.Buffer + termA := slog.NewTextHandler(&bufA, nil) + termB := slog.NewTextHandler(&bufB, nil) + var holder atomic.Pointer[versionedHandler] + holder.Store(&versionedHandler{version: 1, h: termA}) + inj := newSamplingInjector(&holder, termA) + rec := slog.NewRecord(time.Unix(0, 0), slog.LevelInfo, "m", 0) + require.NoError(t, inj.Handle(context.Background(), rec)) + require.Contains(t, bufA.String(), "m") + holder.Store(&versionedHandler{version: 2, h: termB}) // reload + require.NoError(t, inj.Handle(context.Background(), rec)) + require.Contains(t, bufB.String(), "m") // now routed to the new root +} diff --git a/internal/util/metrics.go b/internal/util/metrics.go index 653f97c954a..850b535fe8f 100644 --- a/internal/util/metrics.go +++ b/internal/util/metrics.go @@ -14,16 +14,3 @@ func MustRegisterOrGet(reg prometheus.Registerer, c prometheus.Collector) promet } return c } - -// MustRegisterOrReturnExisting will attempt to register the supplied collector into the register. If it's already -// registered, it will return that one otherwise nil. -// In case that the register procedure fails with something other than already registered, this will panic. -func MustRegisterOrReturnExisting(reg prometheus.Registerer, c prometheus.Collector) prometheus.Collector { - if err := reg.Register(c); err != nil { - if are, ok := err.(prometheus.AlreadyRegisteredError); ok { - return are.ExistingCollector - } - panic(err) - } - return nil -} diff --git a/internal/util/metricsutil/metricsutil.go b/internal/util/metricsutil/metricsutil.go new file mode 100644 index 00000000000..a55932287ae --- /dev/null +++ b/internal/util/metricsutil/metricsutil.go @@ -0,0 +1,21 @@ +// Package metricsutil holds small Prometheus helpers with no dependencies +// beyond the Prometheus client. Packages that cannot import internal/util, +// for example because internal/util already imports them, can still use +// these helpers by depending on this leaf package instead. +package metricsutil + +import "github.com/prometheus/client_golang/prometheus" + +// MustRegisterOrReturnExisting registers c on reg. If c is already +// registered, for example because multiple callers share one registerer, it +// returns the existing collector instead of panicking. If registration succeeds, +// it returns nil. If registration fails for any other reason, it panics. +func MustRegisterOrReturnExisting(reg prometheus.Registerer, c prometheus.Collector) prometheus.Collector { + if err := reg.Register(c); err != nil { + if are, ok := err.(prometheus.AlreadyRegisteredError); ok { + return are.ExistingCollector + } + panic(err) + } + return nil +}