From 8fe0f7e6a27cf3efdf56433639f1b24ee8b542f1 Mon Sep 17 00:00:00 2001 From: Joel Wetzell Date: Sun, 24 May 2026 14:19:58 -0500 Subject: [PATCH] add rate limiting filter --- internal/processor/filter-rate.go | 58 ++++++++++ internal/processor/test/filter-rate_test.go | 116 ++++++++++++++++++++ 2 files changed, 174 insertions(+) create mode 100644 internal/processor/filter-rate.go create mode 100644 internal/processor/test/filter-rate_test.go diff --git a/internal/processor/filter-rate.go b/internal/processor/filter-rate.go new file mode 100644 index 0000000..9e0c6b9 --- /dev/null +++ b/internal/processor/filter-rate.go @@ -0,0 +1,58 @@ +package processor + +import ( + "context" + "fmt" + + "github.com/google/jsonschema-go/jsonschema" + "github.com/jwetzell/showbridge-go/internal/common" + "github.com/jwetzell/showbridge-go/internal/config" + "golang.org/x/time/rate" +) + +func init() { + RegisterProcessor(ProcessorRegistration{ + Type: "filter.rate", + Title: "Filter by Rate", + ParamsSchema: &jsonschema.Schema{ + Type: "object", + Properties: map[string]*jsonschema.Schema{ + "rate": { + Type: "integer", + Title: "Rate", + Description: "The number of events to allow per second.", + }, + }, + Required: []string{"rate"}, + }, + New: func(config config.ProcessorConfig) (Processor, error) { + params := config.Params + rateInt, err := params.GetInt("rate") + if err != nil { + return nil, fmt.Errorf("filter.rate rate error: %w", err) + } + + limiter := rate.NewLimiter(rate.Limit(rateInt), rateInt*2) + + return &FilterRate{config: config, limiter: limiter}, nil + }, + }) +} + +type FilterRate struct { + config config.ProcessorConfig + limiter *rate.Limiter +} + +func (fc *FilterRate) Process(ctx context.Context, wrappedPayload common.WrappedPayload) (common.WrappedPayload, error) { + err := fc.limiter.Wait(ctx) + if err != nil { + wrappedPayload.End = true + return wrappedPayload, err + } + return wrappedPayload, nil +} + +func (fc *FilterRate) Type() string { + return fc.config.Type +} diff --git a/internal/processor/test/filter-rate_test.go b/internal/processor/test/filter-rate_test.go new file mode 100644 index 0000000..0e28fd0 --- /dev/null +++ b/internal/processor/test/filter-rate_test.go @@ -0,0 +1,116 @@ +package processor_test + +import ( + "testing" + + "github.com/jwetzell/showbridge-go/internal/common" + "github.com/jwetzell/showbridge-go/internal/config" + "github.com/jwetzell/showbridge-go/internal/processor" + "github.com/jwetzell/showbridge-go/internal/test" +) + +func TestFilterRateFromRegistry(t *testing.T) { + registration, ok := processor.ProcessorRegistry["filter.rate"] + if !ok { + t.Fatalf("filter.rate processor not registered") + } + + processorInstance, err := registration.New(config.ProcessorConfig{ + Type: "filter.rate", + Params: map[string]any{ + "rate": 1, + }, + }) + if err != nil { + t.Fatalf("failed to create filter.rate processor: %s", err) + } + + if processorInstance.Type() != "filter.rate" { + t.Fatalf("filter.rate processor has wrong type: %s", processorInstance.Type()) + } +} + +func TestGoodFilterRate(t *testing.T) { + testCases := []struct { + name string + params map[string]any + payload any + match bool + }{} + + for _, testCase := range testCases { + t.Run(testCase.name, func(t *testing.T) { + registration, ok := processor.ProcessorRegistry["filter.rate"] + if !ok { + t.Fatalf("filter.rate processor not registered") + } + + processorInstance, err := registration.New(config.ProcessorConfig{ + Type: "filter.rate", + Params: testCase.params, + }) + + if err != nil { + t.Fatalf("filter.rate failed to create processor: %s", err) + } + + _, err = processorInstance.Process(t.Context(), common.WrappedPayload{Payload: testCase.payload}) + // TODO(jwetzell): figure out how to test the rate limiting behavior + }) + } +} + +func TestBadFilterRate(t *testing.T) { + tests := []struct { + name string + params map[string]any + payload any + errorString string + }{ + { + name: "no rate parameter", + params: map[string]any{ + // no rate parameter + }, + payload: test.TestStruct{}, + errorString: "filter.rate rate error: not found", + }, + { + name: "non-int rate parameter", + params: map[string]any{ + "rate": "12345", + }, + payload: test.TestStruct{}, + errorString: "filter.rate rate error: not a number", + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + registration, ok := processor.ProcessorRegistry["filter.rate"] + if !ok { + t.Fatalf("filter.rate processor not registered") + } + + processorInstance, err := registration.New(config.ProcessorConfig{ + Type: "filter.rate", + Params: test.params, + }) + if err != nil { + if err.Error() != test.errorString { + t.Fatalf("filter.rate got error '%s', expected '%s'", err.Error(), test.errorString) + } + return + } + got, err := processorInstance.Process(t.Context(), common.WrappedPayload{Payload: test.payload}) + + if err == nil { + t.Fatalf("filter.rate expected to fail but succeeded, got: %v", got) + + } + if err.Error() != test.errorString { + t.Fatalf("filter.rate got error '%s', expected '%s'", err.Error(), test.errorString) + } + }) + } +}