mirror of
https://github.com/open-policy-agent/opa.git
synced 2026-08-13 03:42:35 -06:00
plugin/decision: set config boundaries to upload_size_limit_bytes (#7563)
Signed-off-by: sspaink <sspaink@styra.com>
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user