From e1c37ba30d6bcd96c10b92598d6ef2fce0617841 Mon Sep 17 00:00:00 2001 From: Sebastian Spaink <3441183+sspaink@users.noreply.github.com> Date: Mon, 19 May 2025 09:09:25 -0500 Subject: [PATCH] plugin/decision: set config boundaries to upload_size_limit_bytes (#7563) Signed-off-by: sspaink --- v1/plugins/discovery/discovery.go | 10 +-- v1/plugins/discovery/discovery_test.go | 10 +-- v1/plugins/logs/plugin.go | 33 +++++++-- v1/plugins/logs/plugin_benchmark_test.go | 6 +- v1/plugins/logs/plugin_test.go | 89 ++++++++++++++++++++++-- 5 files changed, 126 insertions(+), 22 deletions(-) diff --git a/v1/plugins/discovery/discovery.go b/v1/plugins/discovery/discovery.go index 117bb53df5..e34d1f2b68 100644 --- a/v1/plugins/discovery/discovery.go +++ b/v1/plugins/discovery/discovery.go @@ -105,13 +105,15 @@ func New(manager *plugins.Manager, opts ...func(*Discovery)) (*Discovery, error) f(result) } + result.logger = manager.Logger().WithFields(map[string]any{"plugin": Name}) + config, err := NewConfigBuilder().WithBytes(manager.Config.Discovery).WithServices(manager.Services()). WithKeyConfigs(manager.PublicKeys()).Parse() if err != nil { return nil, err } else if config == nil { - if _, err := getPluginSet(result.factories, manager, manager.Config, result.metrics, nil); err != nil { + if _, err := getPluginSet(result.factories, manager, manager.Config, result.metrics, result.logger, nil); err != nil { return nil, err } return result, nil @@ -509,7 +511,7 @@ func (c *Discovery) processBundle(ctx context.Context, b *bundleApi.Bundle) (*pl return nil, err } - ps, err := getPluginSet(c.factories, c.manager, overriddenConfig, c.metrics, c.config.Trigger) + ps, err := getPluginSet(c.factories, c.manager, overriddenConfig, c.metrics, c.logger, c.config.Trigger) if err != nil { return nil, err } @@ -584,7 +586,7 @@ type pluginfactory struct { config any } -func getPluginSet(factories map[string]plugins.Factory, manager *plugins.Manager, config *config.Config, m metrics.Metrics, trigger *plugins.TriggerMode) (*pluginSet, error) { +func getPluginSet(factories map[string]plugins.Factory, manager *plugins.Manager, config *config.Config, m metrics.Metrics, l logging.Logger, trigger *plugins.TriggerMode) (*pluginSet, error) { // Parse and validate plugin configurations. pluginNames := []string{} @@ -628,7 +630,7 @@ func getPluginSet(factories map[string]plugins.Factory, manager *plugins.Manager } decisionLogsConfig, err := logs.NewConfigBuilder().WithBytes(config.DecisionLogs).WithServices(manager.Services()). - WithPlugins(pluginNames).WithTriggerMode(trigger).Parse() + WithPlugins(pluginNames).WithTriggerMode(trigger).WithLogger(l).Parse() if err != nil { return nil, err } diff --git a/v1/plugins/discovery/discovery_test.go b/v1/plugins/discovery/discovery_test.go index 464c68708a..fddc7f509f 100644 --- a/v1/plugins/discovery/discovery_test.go +++ b/v1/plugins/discovery/discovery_test.go @@ -3239,7 +3239,7 @@ bundle: ` manager := getTestManager(t, conf) trigger := plugins.TriggerManual - _, err := getPluginSet(nil, manager, manager.Config, nil, &trigger) + _, err := getPluginSet(nil, manager, manager.Config, nil, nil, &trigger) if err != nil { t.Fatalf("Unexpected error: %s", err) } @@ -3276,7 +3276,7 @@ bundles: ` manager := getTestManager(t, conf) trigger := plugins.TriggerManual - _, err := getPluginSet(nil, manager, manager.Config, nil, &trigger) + _, err := getPluginSet(nil, manager, manager.Config, nil, nil, &trigger) if err != nil { t.Fatalf("Unexpected error: %s", err) } @@ -3336,7 +3336,7 @@ bundles: manager := getTestManager(t, tc.conf) trigger := plugins.TriggerManual - _, err := getPluginSet(nil, manager, manager.Config, nil, &trigger) + _, err := getPluginSet(nil, manager, manager.Config, nil, nil, &trigger) if tc.wantErr { if err == nil { @@ -3401,7 +3401,7 @@ decision_logs: manager := getTestManager(t, tc.conf) trigger := plugins.TriggerManual - _, err := getPluginSet(nil, manager, manager.Config, nil, &trigger) + _, err := getPluginSet(nil, manager, manager.Config, nil, nil, &trigger) if tc.wantErr { if err == nil { @@ -3472,7 +3472,7 @@ status: manager := getTestManager(t, tc.conf) trigger := plugins.TriggerManual - _, err := getPluginSet(nil, manager, manager.Config, nil, &trigger) + _, err := getPluginSet(nil, manager, manager.Config, nil, nil, &trigger) if tc.wantErr { if err == nil { diff --git a/v1/plugins/logs/plugin.go b/v1/plugins/logs/plugin.go index 3d2150c49a..acf197fb42 100644 --- a/v1/plugins/logs/plugin.go +++ b/v1/plugins/logs/plugin.go @@ -253,8 +253,10 @@ const ( defaultMinDelaySeconds = int64(300) defaultMaxDelaySeconds = int64(600) defaultBufferSizeLimitEvents = int64(10000) - defaultUploadSizeLimitBytes = int64(32768) // 32KB limit - defaultBufferSizeLimitBytes = int64(0) // unlimited + defaultUploadSizeLimitBytes = int64(32768) // 32KB limit + minUploadSizeLimitBytes = int64(90) // A single event with a decision ID (69 bytes) + empty gzip file (21 bytes) + maxUploadSizeLimitBytes = int64(4294967296) // about 4GB + defaultBufferSizeLimitBytes = int64(0) // unlimited defaultMaskDecisionPath = "/system/log/mask" defaultDropDecisionPath = "/system/log/drop" logRateLimitExDropCounterName = "decision_logs_dropped_rate_limit_exceeded" @@ -303,7 +305,7 @@ type Config struct { dropDecisionRef ast.Ref } -func (c *Config) validateAndInjectDefaults(services []string, pluginsList []string, trigger *plugins.TriggerMode) error { +func (c *Config) validateAndInjectDefaults(services []string, pluginsList []string, trigger *plugins.TriggerMode, l logging.Logger) error { if c.Plugin != nil { var found bool @@ -372,7 +374,22 @@ func (c *Config) validateAndInjectDefaults(services []string, pluginsList []stri uploadLimit = *c.Reporting.UploadSizeLimitBytes } - c.Reporting.UploadSizeLimitBytes = &uploadLimit + switch { + case uploadLimit > maxUploadSizeLimitBytes: + maxUploadLimit := maxUploadSizeLimitBytes + c.Reporting.UploadSizeLimitBytes = &maxUploadLimit + if l != nil { + l.Warn("the configured `upload_size_limit_bytes` (%d) has been set to the maximum limit (%d)", uploadLimit, maxUploadLimit) + } + case uploadLimit < minUploadSizeLimitBytes: + minUploadLimit := minUploadSizeLimitBytes + c.Reporting.UploadSizeLimitBytes = &minUploadLimit + if l != nil { + l.Warn("the configured `upload_size_limit_bytes` (%d) has been set to the minimum limit (%d)", uploadLimit, minUploadLimit) + } + default: + c.Reporting.UploadSizeLimitBytes = &uploadLimit + } if c.Reporting.BufferType == "" { c.Reporting.BufferType = sizeBufferType @@ -511,6 +528,7 @@ type ConfigBuilder struct { services []string plugins []string trigger *plugins.TriggerMode + logger logging.Logger } // NewConfigBuilder returns a new ConfigBuilder to build and parse the plugin config. @@ -518,6 +536,11 @@ func NewConfigBuilder() *ConfigBuilder { return &ConfigBuilder{} } +func (b *ConfigBuilder) WithLogger(l logging.Logger) *ConfigBuilder { + b.logger = l + return b +} + // WithBytes sets the raw plugin config. func (b *ConfigBuilder) WithBytes(config []byte) *ConfigBuilder { b.raw = config @@ -559,7 +582,7 @@ func (b *ConfigBuilder) Parse() (*Config, error) { return nil, nil } - if err := parsedConfig.validateAndInjectDefaults(b.services, b.plugins, b.trigger); err != nil { + if err := parsedConfig.validateAndInjectDefaults(b.services, b.plugins, b.trigger, b.logger); err != nil { return nil, err } diff --git a/v1/plugins/logs/plugin_benchmark_test.go b/v1/plugins/logs/plugin_benchmark_test.go index 1dfb08e8a9..20dce0e609 100644 --- a/v1/plugins/logs/plugin_benchmark_test.go +++ b/v1/plugins/logs/plugin_benchmark_test.go @@ -149,7 +149,7 @@ func BenchmarkMaskingNop(b *testing.B) { cfg := &Config{Service: "svc"} t := plugins.DefaultTriggerMode - if err := cfg.validateAndInjectDefaults([]string{"svc"}, nil, &t); err != nil { + if err := cfg.validateAndInjectDefaults([]string{"svc"}, nil, &t, nil); err != nil { b.Fatal(err) } plugin := New(cfg, manager) @@ -187,7 +187,7 @@ func BenchmarkMaskingRuleCountsNop(b *testing.B) { cfg := &Config{Service: "svc"} t := plugins.DefaultTriggerMode - if err := cfg.validateAndInjectDefaults([]string{"svc"}, nil, &t); err != nil { + if err := cfg.validateAndInjectDefaults([]string{"svc"}, nil, &t, nil); err != nil { b.Fatal(err) } plugin := New(cfg, manager) @@ -240,7 +240,7 @@ func BenchmarkMaskingErase(b *testing.B) { cfg := &Config{Service: "svc"} t := plugins.DefaultTriggerMode - if err := cfg.validateAndInjectDefaults([]string{"svc"}, nil, &t); err != nil { + if err := cfg.validateAndInjectDefaults([]string{"svc"}, nil, &t, nil); err != nil { b.Fatal(err) } plugin := New(cfg, manager) diff --git a/v1/plugins/logs/plugin_test.go b/v1/plugins/logs/plugin_test.go index f567fc8e97..22ffbe0da8 100644 --- a/v1/plugins/logs/plugin_test.go +++ b/v1/plugins/logs/plugin_test.go @@ -2587,7 +2587,7 @@ func TestPluginMasking(t *testing.T) { // Instantiate the plugin. cfg := &Config{Service: "svc"} trigger := plugins.DefaultTriggerMode - cfg.validateAndInjectDefaults([]string{"svc"}, nil, &trigger) + cfg.validateAndInjectDefaults([]string{"svc"}, nil, &trigger, nil) plugin := New(cfg, manager) @@ -2637,7 +2637,7 @@ func TestPluginMasking(t *testing.T) { // Reconfigure and ensure that mask is invalidated. maskDecision := "dead/beef" newConfig := &Config{Service: "svc", MaskDecision: &maskDecision} - if err := newConfig.validateAndInjectDefaults([]string{"svc"}, nil, &trigger); err != nil { + if err := newConfig.validateAndInjectDefaults([]string{"svc"}, nil, &trigger, nil); err != nil { t.Fatal(err) } @@ -2737,7 +2737,7 @@ func TestPluginDrop(t *testing.T) { // Instantiate the plugin. cfg := &Config{Service: "svc"} trigger := plugins.DefaultTriggerMode - cfg.validateAndInjectDefaults([]string{"svc"}, nil, &trigger) + cfg.validateAndInjectDefaults([]string{"svc"}, nil, &trigger, nil) plugin := New(cfg, manager) @@ -2808,7 +2808,7 @@ func TestPluginMaskErrorHandling(t *testing.T) { // Instantiate the plugin. cfg := &Config{Service: "svc"} trigger := plugins.DefaultTriggerMode - cfg.validateAndInjectDefaults([]string{"svc"}, nil, &trigger) + cfg.validateAndInjectDefaults([]string{"svc"}, nil, &trigger, nil) plugin := New(cfg, manager) @@ -2884,7 +2884,7 @@ func TestPluginDropErrorHandling(t *testing.T) { // Instantiate the plugin. cfg := &Config{Service: "svc"} trigger := plugins.DefaultTriggerMode - cfg.validateAndInjectDefaults([]string{"svc"}, nil, &trigger) + cfg.validateAndInjectDefaults([]string{"svc"}, nil, &trigger, nil) plugin := New(cfg, manager) @@ -3654,3 +3654,82 @@ func testStatus() *bundle.Status { LastSuccessfulActivation: tActivate, } } + +func TestConfigUploadLimit(t *testing.T) { + tests := []struct { + name string + limit int64 + expectedLimit int64 + expectedLog string + expectedErr string + }{ + { + name: "exceed maximum limit", + limit: int64(8589934592), + expectedLimit: maxUploadSizeLimitBytes, + expectedLog: "the configured `upload_size_limit_bytes` (8589934592) has been set to the maximum limit (4294967296)", + }, + { + name: "nothing changes", + limit: 1000, + expectedLimit: 1000, + }, + { + name: "negative limit", + limit: -1, + expectedLimit: minUploadSizeLimitBytes, + expectedLog: "the configured `upload_size_limit_bytes` (-1) has been set to the minimum limit (90)", + }, + { + name: "equal to minimum", + limit: minUploadSizeLimitBytes, + expectedLimit: minUploadSizeLimitBytes, + }, + { + name: "equal to maximum", + limit: maxUploadSizeLimitBytes, + expectedLimit: maxUploadSizeLimitBytes, + }, + } + + for _, tc := range tests { + t.Run(tc.name, func(t *testing.T) { + + testLogger := test.New() + + cfg := &Config{ + Service: "svc", + Reporting: ReportingConfig{ + UploadSizeLimitBytes: &tc.limit, + }, + } + trigger := plugins.DefaultTriggerMode + if err := cfg.validateAndInjectDefaults([]string{"svc"}, nil, &trigger, testLogger); err != nil { + if tc.expectedErr != "" { + if tc.expectedErr != err.Error() { + t.Fatalf("Expected error to be `%s` but got `%s`", tc.expectedErr, err.Error()) + } else { + return + } + } else { + t.Fatal(err) + } + } + + if *cfg.Reporting.UploadSizeLimitBytes != tc.expectedLimit { + t.Fatalf("Expected upload limit to be %d but got %d", tc.expectedLimit, cfg.Reporting.UploadSizeLimitBytes) + } + + if tc.expectedLog != "" { + e := testLogger.Entries() + if e[0].Message != tc.expectedLog { + t.Fatalf("Expected log to be %s but got %s", tc.expectedLog, e[0].Message) + } + } else { + if len(testLogger.Entries()) != 0 { + t.Fatalf("Expected log to be empty but got %s", testLogger.Entries()[0].Message) + } + } + }) + } +}