mirror of
https://github.com/jwetzell/showbridge-go.git
synced 2026-09-10 23:49:25 +00:00
Merge pull request #164 from jwetzell/router-output-rename
rename router.output to module.output
This commit is contained in:
@@ -11,37 +11,53 @@ import (
|
|||||||
"github.com/jwetzell/showbridge-go/internal/config"
|
"github.com/jwetzell/showbridge-go/internal/config"
|
||||||
)
|
)
|
||||||
|
|
||||||
type RouterOutput struct {
|
type ModuleOutput struct {
|
||||||
config config.ProcessorConfig
|
config config.ProcessorConfig
|
||||||
ModuleId string
|
ModuleId string
|
||||||
logger *slog.Logger
|
logger *slog.Logger
|
||||||
|
module common.OutputModule
|
||||||
}
|
}
|
||||||
|
|
||||||
func (ro *RouterOutput) Process(ctx context.Context, wrappedPayload common.WrappedPayload) (common.WrappedPayload, error) {
|
func (ro *ModuleOutput) Process(ctx context.Context, wrappedPayload common.WrappedPayload) (common.WrappedPayload, error) {
|
||||||
|
|
||||||
if wrappedPayload.Router == nil {
|
if ro.module == nil {
|
||||||
wrappedPayload.End = true
|
if wrappedPayload.Modules == nil {
|
||||||
return wrappedPayload, errors.New("router.output no router found")
|
wrappedPayload.End = true
|
||||||
|
return wrappedPayload, errors.New("module.output wrapped payload has no modules")
|
||||||
|
}
|
||||||
|
|
||||||
|
module, ok := wrappedPayload.Modules[ro.ModuleId]
|
||||||
|
if !ok {
|
||||||
|
wrappedPayload.End = true
|
||||||
|
return wrappedPayload, fmt.Errorf("module.output unable to find module with id: %s", ro.ModuleId)
|
||||||
|
}
|
||||||
|
|
||||||
|
outputModule, ok := module.(common.OutputModule)
|
||||||
|
if !ok {
|
||||||
|
wrappedPayload.End = true
|
||||||
|
return wrappedPayload, fmt.Errorf("module.output module with id %s is not an OutputModule", ro.ModuleId)
|
||||||
|
}
|
||||||
|
ro.module = outputModule
|
||||||
}
|
}
|
||||||
|
|
||||||
err := wrappedPayload.Router.HandleOutput(ctx, ro.ModuleId, wrappedPayload.Payload)
|
err := ro.module.Output(ctx, wrappedPayload.Payload)
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
wrappedPayload.End = true
|
wrappedPayload.End = true
|
||||||
return wrappedPayload, fmt.Errorf("router.output failed to send output: %w", err)
|
return wrappedPayload, fmt.Errorf("module.output failed to send output: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
return wrappedPayload, nil
|
return wrappedPayload, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func (ro *RouterOutput) Type() string {
|
func (ro *ModuleOutput) Type() string {
|
||||||
return ro.config.Type
|
return ro.config.Type
|
||||||
}
|
}
|
||||||
|
|
||||||
func init() {
|
func init() {
|
||||||
RegisterProcessor(ProcessorRegistration{
|
RegisterProcessor(ProcessorRegistration{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Title: "Router Output",
|
Title: "Module Output",
|
||||||
ParamsSchema: &jsonschema.Schema{
|
ParamsSchema: &jsonschema.Schema{
|
||||||
Type: "object",
|
Type: "object",
|
||||||
Properties: map[string]*jsonschema.Schema{
|
Properties: map[string]*jsonschema.Schema{
|
||||||
@@ -61,10 +77,10 @@ func init() {
|
|||||||
moduleId, err := params.GetString("module")
|
moduleId, err := params.GetString("module")
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("router.output module error: %w", err)
|
return nil, fmt.Errorf("module.output module error: %w", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
return &RouterOutput{config: config, ModuleId: moduleId, logger: slog.Default().With("component", "processor", "type", config.Type)}, nil
|
return &ModuleOutput{config: config, ModuleId: moduleId, logger: slog.Default().With("component", "processor", "type", config.Type)}, nil
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -0,0 +1,155 @@
|
|||||||
|
package processor_test
|
||||||
|
|
||||||
|
import (
|
||||||
|
"reflect"
|
||||||
|
"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 TestModuleOutputFromRegistry(t *testing.T) {
|
||||||
|
registration, ok := processor.ProcessorRegistry["module.output"]
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("module.output processor not registered")
|
||||||
|
}
|
||||||
|
|
||||||
|
processorInstance, err := registration.New(config.ProcessorConfig{
|
||||||
|
Type: "module.output",
|
||||||
|
Params: config.Params{
|
||||||
|
"module": "test",
|
||||||
|
},
|
||||||
|
})
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("failed to create module.output processor: %s", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if processorInstance.Type() != "module.output" {
|
||||||
|
t.Fatalf("module.output processor has wrong type: %s", processorInstance.Type())
|
||||||
|
}
|
||||||
|
|
||||||
|
payload := "test"
|
||||||
|
expected := "test"
|
||||||
|
|
||||||
|
got, err := processorInstance.Process(t.Context(), common.WrappedPayload{
|
||||||
|
Router: test.GetNewTestRouter(),
|
||||||
|
Modules: map[string]common.Module{"test": &test.TestOutputModule{}},
|
||||||
|
Payload: payload,
|
||||||
|
})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("module.output processing failed: %s", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if got.Payload != expected {
|
||||||
|
t.Fatalf("module.output got %+v, expected %+v", got, expected)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestGoodModuleOutput(t *testing.T) {
|
||||||
|
|
||||||
|
testCases := []struct {
|
||||||
|
name string
|
||||||
|
params map[string]any
|
||||||
|
payload any
|
||||||
|
expected any
|
||||||
|
}{}
|
||||||
|
|
||||||
|
for _, testCase := range testCases {
|
||||||
|
t.Run(testCase.name, func(t *testing.T) {
|
||||||
|
|
||||||
|
registration, ok := processor.ProcessorRegistry["module.output"]
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("module.output processor not registered")
|
||||||
|
}
|
||||||
|
|
||||||
|
processorInstance, err := registration.New(config.ProcessorConfig{
|
||||||
|
Type: "module.output",
|
||||||
|
Params: testCase.params,
|
||||||
|
})
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("module.output failed to create processor: %s", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
got, err := processorInstance.Process(t.Context(), common.WrappedPayload{Payload: testCase.payload})
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("module.output processing failed: %s", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
if !reflect.DeepEqual(got.Payload, testCase.expected) {
|
||||||
|
t.Fatalf("module.output got %+v (%T), expected %+v (%T)", got.Payload, got.Payload, testCase.expected, testCase.expected)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestBadModuleOutput(t *testing.T) {
|
||||||
|
testCases := []struct {
|
||||||
|
name string
|
||||||
|
params map[string]any
|
||||||
|
payload any
|
||||||
|
modules map[string]common.Module
|
||||||
|
errorString string
|
||||||
|
}{
|
||||||
|
{
|
||||||
|
name: "no module param",
|
||||||
|
params: map[string]any{},
|
||||||
|
payload: "test",
|
||||||
|
modules: map[string]common.Module{"test": &test.TestModule{}},
|
||||||
|
errorString: "module.output module error: not found",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "non-string module",
|
||||||
|
params: map[string]any{
|
||||||
|
"module": 123,
|
||||||
|
},
|
||||||
|
payload: "test",
|
||||||
|
modules: map[string]common.Module{"test": &test.TestModule{}},
|
||||||
|
errorString: "module.output module error: not a string",
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "modules not found in context",
|
||||||
|
params: map[string]any{
|
||||||
|
"module": "test",
|
||||||
|
},
|
||||||
|
payload: "test",
|
||||||
|
modules: nil,
|
||||||
|
errorString: "module.output wrapped payload has no modules",
|
||||||
|
},
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, testCase := range testCases {
|
||||||
|
t.Run(testCase.name, func(t *testing.T) {
|
||||||
|
|
||||||
|
registration, ok := processor.ProcessorRegistry["module.output"]
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("module.output processor not registered")
|
||||||
|
}
|
||||||
|
|
||||||
|
processorInstance, err := registration.New(config.ProcessorConfig{
|
||||||
|
Type: "module.output",
|
||||||
|
Params: testCase.params,
|
||||||
|
})
|
||||||
|
|
||||||
|
if err != nil {
|
||||||
|
if testCase.errorString != err.Error() {
|
||||||
|
t.Fatalf("module.output got error '%s', expected '%s'", err.Error(), testCase.errorString)
|
||||||
|
}
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
got, err := processorInstance.Process(t.Context(), common.WrappedPayload{Modules: testCase.modules, Payload: testCase.payload})
|
||||||
|
|
||||||
|
if err == nil {
|
||||||
|
t.Fatalf("module.output expected to fail but succeeded, got: %v", got)
|
||||||
|
}
|
||||||
|
|
||||||
|
if err.Error() != testCase.errorString {
|
||||||
|
t.Fatalf("module.output got error '%s', expected '%s'", err.Error(), testCase.errorString)
|
||||||
|
}
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -10,25 +10,25 @@ import (
|
|||||||
"github.com/jwetzell/showbridge-go/internal/test"
|
"github.com/jwetzell/showbridge-go/internal/test"
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestRouterOutputFromRegistry(t *testing.T) {
|
func TestRouterInputFromRegistry(t *testing.T) {
|
||||||
registration, ok := processor.ProcessorRegistry["router.output"]
|
registration, ok := processor.ProcessorRegistry["router.input"]
|
||||||
if !ok {
|
if !ok {
|
||||||
t.Fatalf("router.output processor not registered")
|
t.Fatalf("router.input processor not registered")
|
||||||
}
|
}
|
||||||
|
|
||||||
processorInstance, err := registration.New(config.ProcessorConfig{
|
processorInstance, err := registration.New(config.ProcessorConfig{
|
||||||
Type: "router.output",
|
Type: "router.input",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "test",
|
"source": "test",
|
||||||
},
|
},
|
||||||
})
|
})
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("failed to create router.output processor: %s", err)
|
t.Fatalf("failed to create router.input processor: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if processorInstance.Type() != "router.output" {
|
if processorInstance.Type() != "router.input" {
|
||||||
t.Fatalf("router.output processor has wrong type: %s", processorInstance.Type())
|
t.Fatalf("router.input processor has wrong type: %s", processorInstance.Type())
|
||||||
}
|
}
|
||||||
|
|
||||||
payload := "test"
|
payload := "test"
|
||||||
@@ -39,15 +39,15 @@ func TestRouterOutputFromRegistry(t *testing.T) {
|
|||||||
Payload: payload,
|
Payload: payload,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("router.output processing failed: %s", err)
|
t.Fatalf("router.input processing failed: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if got.Payload != expected {
|
if got.Payload != expected {
|
||||||
t.Fatalf("router.output got %+v, expected %+v", got, expected)
|
t.Fatalf("router.input got %+v, expected %+v", got, expected)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestGoodRouterOutput(t *testing.T) {
|
func TestGoodRouterInput(t *testing.T) {
|
||||||
|
|
||||||
testCases := []struct {
|
testCases := []struct {
|
||||||
name string
|
name string
|
||||||
@@ -59,33 +59,33 @@ func TestGoodRouterOutput(t *testing.T) {
|
|||||||
for _, testCase := range testCases {
|
for _, testCase := range testCases {
|
||||||
t.Run(testCase.name, func(t *testing.T) {
|
t.Run(testCase.name, func(t *testing.T) {
|
||||||
|
|
||||||
registration, ok := processor.ProcessorRegistry["router.output"]
|
registration, ok := processor.ProcessorRegistry["router.input"]
|
||||||
if !ok {
|
if !ok {
|
||||||
t.Fatalf("router.output processor not registered")
|
t.Fatalf("router.input processor not registered")
|
||||||
}
|
}
|
||||||
|
|
||||||
processorInstance, err := registration.New(config.ProcessorConfig{
|
processorInstance, err := registration.New(config.ProcessorConfig{
|
||||||
Type: "router.output",
|
Type: "router.input",
|
||||||
Params: testCase.params,
|
Params: testCase.params,
|
||||||
})
|
})
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("router.output failed to create processor: %s", err)
|
t.Fatalf("router.input failed to create processor: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
got, err := processorInstance.Process(t.Context(), common.WrappedPayload{Payload: testCase.payload})
|
got, err := processorInstance.Process(t.Context(), common.WrappedPayload{Payload: testCase.payload})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
t.Fatalf("router.output processing failed: %s", err)
|
t.Fatalf("router.input processing failed: %s", err)
|
||||||
}
|
}
|
||||||
|
|
||||||
if !reflect.DeepEqual(got.Payload, testCase.expected) {
|
if !reflect.DeepEqual(got.Payload, testCase.expected) {
|
||||||
t.Fatalf("router.output got %+v (%T), expected %+v (%T)", got.Payload, got.Payload, testCase.expected, testCase.expected)
|
t.Fatalf("router.input got %+v (%T), expected %+v (%T)", got.Payload, got.Payload, testCase.expected, testCase.expected)
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestBadRouterOutput(t *testing.T) {
|
func TestBadRouterInput(t *testing.T) {
|
||||||
testCases := []struct {
|
testCases := []struct {
|
||||||
name string
|
name string
|
||||||
params map[string]any
|
params map[string]any
|
||||||
@@ -94,49 +94,48 @@ func TestBadRouterOutput(t *testing.T) {
|
|||||||
errorString string
|
errorString string
|
||||||
}{
|
}{
|
||||||
{
|
{
|
||||||
name: "no module param",
|
name: "no source param",
|
||||||
params: map[string]any{},
|
params: map[string]any{},
|
||||||
payload: "test",
|
payload: "test",
|
||||||
router: test.GetNewTestRouter(),
|
router: test.GetNewTestRouter(),
|
||||||
errorString: "router.output module error: not found",
|
errorString: "router.input source error: not found",
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
name: "non-string module",
|
name: "non-string source",
|
||||||
params: map[string]any{
|
params: map[string]any{
|
||||||
"module": 123,
|
"source": 123,
|
||||||
},
|
},
|
||||||
payload: "test",
|
payload: "test",
|
||||||
router: test.GetNewTestRouter(),
|
router: test.GetNewTestRouter(),
|
||||||
|
errorString: "router.input source error: not a string",
|
||||||
errorString: "router.output module error: not a string",
|
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
name: "router not found in context",
|
name: "router not found in context",
|
||||||
params: map[string]any{
|
params: map[string]any{
|
||||||
"module": "test",
|
"source": "test",
|
||||||
},
|
},
|
||||||
payload: "test",
|
payload: "test",
|
||||||
router: nil,
|
router: nil,
|
||||||
errorString: "router.output no router found",
|
errorString: "router.input no router found",
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
for _, testCase := range testCases {
|
for _, testCase := range testCases {
|
||||||
t.Run(testCase.name, func(t *testing.T) {
|
t.Run(testCase.name, func(t *testing.T) {
|
||||||
|
|
||||||
registration, ok := processor.ProcessorRegistry["router.output"]
|
registration, ok := processor.ProcessorRegistry["router.input"]
|
||||||
if !ok {
|
if !ok {
|
||||||
t.Fatalf("router.output processor not registered")
|
t.Fatalf("router.input processor not registered")
|
||||||
}
|
}
|
||||||
|
|
||||||
processorInstance, err := registration.New(config.ProcessorConfig{
|
processorInstance, err := registration.New(config.ProcessorConfig{
|
||||||
Type: "router.output",
|
Type: "router.input",
|
||||||
Params: testCase.params,
|
Params: testCase.params,
|
||||||
})
|
})
|
||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
if testCase.errorString != err.Error() {
|
if testCase.errorString != err.Error() {
|
||||||
t.Fatalf("router.output got error '%s', expected '%s'", err.Error(), testCase.errorString)
|
t.Fatalf("router.input got error '%s', expected '%s'", err.Error(), testCase.errorString)
|
||||||
}
|
}
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
@@ -144,11 +143,11 @@ func TestBadRouterOutput(t *testing.T) {
|
|||||||
got, err := processorInstance.Process(t.Context(), common.WrappedPayload{Router: testCase.router, Payload: testCase.payload})
|
got, err := processorInstance.Process(t.Context(), common.WrappedPayload{Router: testCase.router, Payload: testCase.payload})
|
||||||
|
|
||||||
if err == nil {
|
if err == nil {
|
||||||
t.Fatalf("router.output expected to fail but succeeded, got: %v", got)
|
t.Fatalf("router.input expected to fail but succeeded, got: %v", got)
|
||||||
}
|
}
|
||||||
|
|
||||||
if err.Error() != testCase.errorString {
|
if err.Error() != testCase.errorString {
|
||||||
t.Fatalf("router.output got error '%s', expected '%s'", err.Error(), testCase.errorString)
|
t.Fatalf("router.input got error '%s', expected '%s'", err.Error(), testCase.errorString)
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,154 +0,0 @@
|
|||||||
package processor_test
|
|
||||||
|
|
||||||
import (
|
|
||||||
"reflect"
|
|
||||||
"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 TestRouterInputFromRegistry(t *testing.T) {
|
|
||||||
registration, ok := processor.ProcessorRegistry["router.input"]
|
|
||||||
if !ok {
|
|
||||||
t.Fatalf("router.input processor not registered")
|
|
||||||
}
|
|
||||||
|
|
||||||
processorInstance, err := registration.New(config.ProcessorConfig{
|
|
||||||
Type: "router.input",
|
|
||||||
Params: config.Params{
|
|
||||||
"source": "test",
|
|
||||||
},
|
|
||||||
})
|
|
||||||
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("failed to create router.input processor: %s", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if processorInstance.Type() != "router.input" {
|
|
||||||
t.Fatalf("router.input processor has wrong type: %s", processorInstance.Type())
|
|
||||||
}
|
|
||||||
|
|
||||||
payload := "test"
|
|
||||||
expected := "test"
|
|
||||||
|
|
||||||
got, err := processorInstance.Process(t.Context(), common.WrappedPayload{
|
|
||||||
Router: test.GetNewTestRouter(),
|
|
||||||
Payload: payload,
|
|
||||||
})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("router.input processing failed: %s", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if got.Payload != expected {
|
|
||||||
t.Fatalf("router.input got %+v, expected %+v", got, expected)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestGoodRouterInput(t *testing.T) {
|
|
||||||
|
|
||||||
testCases := []struct {
|
|
||||||
name string
|
|
||||||
params map[string]any
|
|
||||||
payload any
|
|
||||||
expected any
|
|
||||||
}{}
|
|
||||||
|
|
||||||
for _, testCase := range testCases {
|
|
||||||
t.Run(testCase.name, func(t *testing.T) {
|
|
||||||
|
|
||||||
registration, ok := processor.ProcessorRegistry["router.input"]
|
|
||||||
if !ok {
|
|
||||||
t.Fatalf("router.input processor not registered")
|
|
||||||
}
|
|
||||||
|
|
||||||
processorInstance, err := registration.New(config.ProcessorConfig{
|
|
||||||
Type: "router.input",
|
|
||||||
Params: testCase.params,
|
|
||||||
})
|
|
||||||
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("router.input failed to create processor: %s", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
got, err := processorInstance.Process(t.Context(), common.WrappedPayload{Payload: testCase.payload})
|
|
||||||
if err != nil {
|
|
||||||
t.Fatalf("router.input processing failed: %s", err)
|
|
||||||
}
|
|
||||||
|
|
||||||
if !reflect.DeepEqual(got.Payload, testCase.expected) {
|
|
||||||
t.Fatalf("router.input got %+v (%T), expected %+v (%T)", got.Payload, got.Payload, testCase.expected, testCase.expected)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
func TestBadRouterInput(t *testing.T) {
|
|
||||||
testCases := []struct {
|
|
||||||
name string
|
|
||||||
params map[string]any
|
|
||||||
payload any
|
|
||||||
router common.RouteIO
|
|
||||||
errorString string
|
|
||||||
}{
|
|
||||||
{
|
|
||||||
name: "no source param",
|
|
||||||
params: map[string]any{},
|
|
||||||
payload: "test",
|
|
||||||
router: test.GetNewTestRouter(),
|
|
||||||
errorString: "router.input source error: not found",
|
|
||||||
},
|
|
||||||
{
|
|
||||||
name: "non-string source",
|
|
||||||
params: map[string]any{
|
|
||||||
"source": 123,
|
|
||||||
},
|
|
||||||
payload: "test",
|
|
||||||
router: test.GetNewTestRouter(),
|
|
||||||
errorString: "router.input source error: not a string",
|
|
||||||
},
|
|
||||||
{
|
|
||||||
name: "router not found in context",
|
|
||||||
params: map[string]any{
|
|
||||||
"source": "test",
|
|
||||||
},
|
|
||||||
payload: "test",
|
|
||||||
router: nil,
|
|
||||||
errorString: "router.input no router found",
|
|
||||||
},
|
|
||||||
}
|
|
||||||
|
|
||||||
for _, testCase := range testCases {
|
|
||||||
t.Run(testCase.name, func(t *testing.T) {
|
|
||||||
|
|
||||||
registration, ok := processor.ProcessorRegistry["router.input"]
|
|
||||||
if !ok {
|
|
||||||
t.Fatalf("router.input processor not registered")
|
|
||||||
}
|
|
||||||
|
|
||||||
processorInstance, err := registration.New(config.ProcessorConfig{
|
|
||||||
Type: "router.input",
|
|
||||||
Params: testCase.params,
|
|
||||||
})
|
|
||||||
|
|
||||||
if err != nil {
|
|
||||||
if testCase.errorString != err.Error() {
|
|
||||||
t.Fatalf("router.input got error '%s', expected '%s'", err.Error(), testCase.errorString)
|
|
||||||
}
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
got, err := processorInstance.Process(t.Context(), common.WrappedPayload{Router: testCase.router, Payload: testCase.payload})
|
|
||||||
|
|
||||||
if err == nil {
|
|
||||||
t.Fatalf("router.input expected to fail but succeeded, got: %v", got)
|
|
||||||
}
|
|
||||||
|
|
||||||
if err.Error() != testCase.errorString {
|
|
||||||
t.Fatalf("router.input got error '%s', expected '%s'", err.Error(), testCase.errorString)
|
|
||||||
}
|
|
||||||
})
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
"github.com/jwetzell/showbridge-go/internal/common"
|
"github.com/jwetzell/showbridge-go/internal/common"
|
||||||
"github.com/jwetzell/showbridge-go/internal/config"
|
"github.com/jwetzell/showbridge-go/internal/config"
|
||||||
"github.com/jwetzell/showbridge-go/internal/route"
|
"github.com/jwetzell/showbridge-go/internal/route"
|
||||||
|
"github.com/jwetzell/showbridge-go/internal/test"
|
||||||
)
|
)
|
||||||
|
|
||||||
func TestRouteCreate(t *testing.T) {
|
func TestRouteCreate(t *testing.T) {
|
||||||
@@ -41,7 +42,7 @@ func TestGoodRouteHandleInput(t *testing.T) {
|
|||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{Type: "string.encode"},
|
{Type: "string.encode"},
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "output",
|
"module": "output",
|
||||||
},
|
},
|
||||||
@@ -57,6 +58,7 @@ func TestGoodRouteHandleInput(t *testing.T) {
|
|||||||
inputData := "test input data"
|
inputData := "test input data"
|
||||||
payload, err := testRoute.ProcessPayload(t.Context(), common.WrappedPayload{
|
payload, err := testRoute.ProcessPayload(t.Context(), common.WrappedPayload{
|
||||||
Router: &MockRouter{},
|
Router: &MockRouter{},
|
||||||
|
Modules: map[string]common.Module{"output": &test.TestOutputModule{}},
|
||||||
Payload: inputData,
|
Payload: inputData,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -79,7 +81,7 @@ func TestRouteHandleInputWithProcessorError(t *testing.T) {
|
|||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{Type: "string.create", Params: map[string]any{"template": "{{.invalid}}}"}},
|
{Type: "string.create", Params: map[string]any{"template": "{{.invalid}}}"}},
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "output",
|
"module": "output",
|
||||||
},
|
},
|
||||||
@@ -95,6 +97,7 @@ func TestRouteHandleInputWithProcessorError(t *testing.T) {
|
|||||||
inputData := "test input data"
|
inputData := "test input data"
|
||||||
_, err = testRoute.ProcessPayload(t.Context(), common.WrappedPayload{
|
_, err = testRoute.ProcessPayload(t.Context(), common.WrappedPayload{
|
||||||
Router: &MockRouter{},
|
Router: &MockRouter{},
|
||||||
|
Modules: map[string]common.Module{"output": &test.TestOutputModule{}},
|
||||||
Payload: inputData,
|
Payload: inputData,
|
||||||
})
|
})
|
||||||
if err == nil {
|
if err == nil {
|
||||||
@@ -107,7 +110,7 @@ func TestRouteHandleNilPayload(t *testing.T) {
|
|||||||
Input: "input",
|
Input: "input",
|
||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "output",
|
"module": "output",
|
||||||
},
|
},
|
||||||
@@ -123,6 +126,7 @@ func TestRouteHandleNilPayload(t *testing.T) {
|
|||||||
|
|
||||||
payload, err := testRoute.ProcessPayload(t.Context(), common.WrappedPayload{
|
payload, err := testRoute.ProcessPayload(t.Context(), common.WrappedPayload{
|
||||||
Router: &MockRouter{},
|
Router: &MockRouter{},
|
||||||
|
Modules: map[string]common.Module{"output": &test.TestOutputModule{}},
|
||||||
Payload: nil,
|
Payload: nil,
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -139,7 +143,7 @@ func TestRouteHandleNilPayloadFromProcessor(t *testing.T) {
|
|||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{Type: "script.js", Params: map[string]any{"program": "payload = undefined"}},
|
{Type: "script.js", Params: map[string]any{"program": "payload = undefined"}},
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "output",
|
"module": "output",
|
||||||
},
|
},
|
||||||
@@ -154,6 +158,7 @@ func TestRouteHandleNilPayloadFromProcessor(t *testing.T) {
|
|||||||
|
|
||||||
_, err = testRoute.ProcessPayload(t.Context(), common.WrappedPayload{
|
_, err = testRoute.ProcessPayload(t.Context(), common.WrappedPayload{
|
||||||
Router: &MockRouter{},
|
Router: &MockRouter{},
|
||||||
|
Modules: map[string]common.Module{"output": &test.TestOutputModule{}},
|
||||||
Payload: "test",
|
Payload: "test",
|
||||||
})
|
})
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
@@ -26,6 +26,35 @@ func (m *TestModule) Id() string {
|
|||||||
return "test"
|
return "test"
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func NewTestOutputModule(id string) *TestOutputModule {
|
||||||
|
return &TestOutputModule{
|
||||||
|
id: id,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
type TestOutputModule struct {
|
||||||
|
id string
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *TestOutputModule) Start(ctx context.Context, router common.RouteIO) error {
|
||||||
|
<-ctx.Done()
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *TestOutputModule) Output(ctx context.Context, payload any) error {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *TestOutputModule) Stop() {}
|
||||||
|
|
||||||
|
func (m *TestOutputModule) Type() string {
|
||||||
|
return "test.output"
|
||||||
|
}
|
||||||
|
|
||||||
|
func (m *TestOutputModule) Id() string {
|
||||||
|
return m.id
|
||||||
|
}
|
||||||
|
|
||||||
func NewTestKVModule(id string) *TestKVModule {
|
func NewTestKVModule(id string) *TestKVModule {
|
||||||
return &TestKVModule{
|
return &TestKVModule{
|
||||||
id: id,
|
id: id,
|
||||||
@@ -63,6 +92,7 @@ func (m *TestKVModule) Set(key string, value any) error {
|
|||||||
m.kvData[key] = value
|
m.kvData[key] = value
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewTestDBModule(id string) *TestDBModule {
|
func NewTestDBModule(id string) *TestDBModule {
|
||||||
return &TestDBModule{
|
return &TestDBModule{
|
||||||
id: id,
|
id: id,
|
||||||
|
|||||||
+9
-9
@@ -203,7 +203,7 @@ func TestRouterInputUnknownDestinationModule(t *testing.T) {
|
|||||||
Input: "mock",
|
Input: "mock",
|
||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "test",
|
"module": "test",
|
||||||
},
|
},
|
||||||
@@ -239,7 +239,7 @@ func TestRouterInputUnknownDestinationModule(t *testing.T) {
|
|||||||
t.Fatalf("router should have returned exactly 1 routing error, got: %d", len(routingErrors))
|
t.Fatalf("router should have returned exactly 1 routing error, got: %d", len(routingErrors))
|
||||||
}
|
}
|
||||||
|
|
||||||
if routingErrors[0].ProcessError.Error() != "router.output failed to send output: no module found for destination id" {
|
if routingErrors[0].ProcessError.Error() != "module.output unable to find module with id: test" {
|
||||||
t.Fatalf("routing output error did not match expected, got: %s", routingErrors[0].ProcessError.Error())
|
t.Fatalf("routing output error did not match expected, got: %s", routingErrors[0].ProcessError.Error())
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -257,7 +257,7 @@ func TestRouterInputNoMatchingRoute(t *testing.T) {
|
|||||||
Input: "test",
|
Input: "test",
|
||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "mock",
|
"module": "mock",
|
||||||
},
|
},
|
||||||
@@ -303,7 +303,7 @@ func TestRouterInputSingleRoute(t *testing.T) {
|
|||||||
Input: "mock",
|
Input: "mock",
|
||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "mock",
|
"module": "mock",
|
||||||
},
|
},
|
||||||
@@ -369,7 +369,7 @@ func TestRouterInputMultipleRoutes(t *testing.T) {
|
|||||||
Input: "mock",
|
Input: "mock",
|
||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "mock",
|
"module": "mock",
|
||||||
},
|
},
|
||||||
@@ -380,7 +380,7 @@ func TestRouterInputMultipleRoutes(t *testing.T) {
|
|||||||
Input: "mock",
|
Input: "mock",
|
||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "mock",
|
"module": "mock",
|
||||||
},
|
},
|
||||||
@@ -391,7 +391,7 @@ func TestRouterInputMultipleRoutes(t *testing.T) {
|
|||||||
Input: "mock",
|
Input: "mock",
|
||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "mock",
|
"module": "mock",
|
||||||
},
|
},
|
||||||
@@ -461,7 +461,7 @@ func TestRouterInputMultipleModules(t *testing.T) {
|
|||||||
Input: "mock1",
|
Input: "mock1",
|
||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "mock1",
|
"module": "mock1",
|
||||||
},
|
},
|
||||||
@@ -472,7 +472,7 @@ func TestRouterInputMultipleModules(t *testing.T) {
|
|||||||
Input: "mock2",
|
Input: "mock2",
|
||||||
Processors: []config.ProcessorConfig{
|
Processors: []config.ProcessorConfig{
|
||||||
{
|
{
|
||||||
Type: "router.output",
|
Type: "module.output",
|
||||||
Params: config.Params{
|
Params: config.Params{
|
||||||
"module": "mock2",
|
"module": "mock2",
|
||||||
},
|
},
|
||||||
|
|||||||
Reference in New Issue
Block a user