diff --git a/go.mod b/go.mod index ad4050dd..c326b52a 100644 --- a/go.mod +++ b/go.mod @@ -65,7 +65,6 @@ require ( github.com/stretchr/testify v1.5.1 github.com/ulikunitz/xz v0.0.0-20180703112113-636d36a76670 // indirect github.com/wcharczuk/go-chart v2.0.1+incompatible - github.com/xdg/scram v0.0.0-20180814205039-7eeb5667e42c go.opencensus.io v0.22.0 go.uber.org/atomic v1.3.2 // indirect go.uber.org/multierr v1.1.0 // indirect @@ -83,4 +82,4 @@ require ( k8s.io/apimachinery v0.0.0-20190612205821-1799e75a0719 k8s.io/client-go v12.0.0+incompatible k8s.io/klog v0.3.3 -) \ No newline at end of file +) diff --git a/pkg/apis/core/v1/validation.go b/pkg/apis/core/v1/validation.go index 98a76ac0..94ac91d7 100644 --- a/pkg/apis/core/v1/validation.go +++ b/pkg/apis/core/v1/validation.go @@ -494,14 +494,14 @@ func (spec MessageQueueTriggerSpec) Validate() error { result = multierror.Append(result, spec.FunctionReference.Validate()) - if !validator.IsValidMessageQueue((string)(spec.MessageQueueType)) { + if !validator.IsValidMessageQueue((string)(spec.MessageQueueType), spec.MqtKind) { result = multierror.Append(result, MakeValidationErr(ErrorUnsupportedType, "MessageQueueTriggerSpec.MessageQueueType", spec.MessageQueueType, "not a supported message queue type")) } else { - if !validator.IsValidTopic((string)(spec.MessageQueueType), spec.Topic) { + if !validator.IsValidTopic((string)(spec.MessageQueueType), spec.Topic, spec.MqtKind) { result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "MessageQueueTriggerSpec.Topic", spec.Topic, "not a valid topic")) } - if len(spec.ResponseTopic) > 0 && !validator.IsValidTopic((string)(spec.MessageQueueType), spec.ResponseTopic) { + if len(spec.ResponseTopic) > 0 && !validator.IsValidTopic((string)(spec.MessageQueueType), spec.ResponseTopic, spec.MqtKind) { result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "MessageQueueTriggerSpec.ResponseTopic", spec.ResponseTopic, "not a valid topic")) } } diff --git a/pkg/fission-cli/cmd/mqtrigger/create.go b/pkg/fission-cli/cmd/mqtrigger/create.go index 1c651e99..d61c07b7 100644 --- a/pkg/fission-cli/cmd/mqtrigger/create.go +++ b/pkg/fission-cli/cmd/mqtrigger/create.go @@ -59,8 +59,10 @@ func (opts *CreateSubCommand) complete(input cli.Input) error { fnName := input.String(flagkey.MqtFnName) fnNamespace := input.String(flagkey.NamespaceFunction) + mqtKind := input.String(flagkey.MqtKind) + mqType := (fv1.MessageQueueType)(input.String(flagkey.MqtMQType)) - if !validator.IsValidMessageQueue((string)(mqType)) { + if !validator.IsValidMessageQueue((string)(mqType), mqtKind) { return errors.New("Unsupported message queue type") } @@ -88,7 +90,7 @@ func (opts *CreateSubCommand) complete(input cli.Input) error { contentType = "application/json" } - err := checkMQTopicAvailability(mqType, topic, respTopic) + err := checkMQTopicAvailability(mqType, mqtKind, topic, respTopic) if err != nil { return err } @@ -119,8 +121,6 @@ func (opts *CreateSubCommand) complete(input cli.Input) error { secret := input.String(flagkey.MqtSecret) - mqtKind := input.String(flagkey.MqtKind) - if input.Bool(flagkey.SpecSave) { specDir := util.GetSpecDir(input) fr, err := spec.ReadSpecs(specDir) @@ -197,9 +197,9 @@ func (opts *CreateSubCommand) run(input cli.Input) error { return nil } -func checkMQTopicAvailability(mqType fv1.MessageQueueType, topics ...string) error { +func checkMQTopicAvailability(mqType fv1.MessageQueueType, mqtKind string, topics ...string) error { for _, t := range topics { - if len(t) > 0 && !validator.IsValidTopic((string)(mqType), t) { + if len(t) > 0 && !validator.IsValidTopic((string)(mqType), t, mqtKind) { return errors.Errorf("invalid topic for %s: %s", mqType, t) } } diff --git a/pkg/mqtrigger/kafka/Dockerfile b/pkg/mqtrigger/kafka/Dockerfile deleted file mode 100644 index 209f4f54..00000000 --- a/pkg/mqtrigger/kafka/Dockerfile +++ /dev/null @@ -1,19 +0,0 @@ -FROM golang:1.12-alpine as builder - -RUN apk add bash ca-certificates git gcc g++ libc-dev - -ARG GOPKG=github.com/fission/fission - -ENV GO111MODULE=on - -WORKDIR /go/src/${GOPKG} -COPY ./ ./ - -WORKDIR /go/src/${GOPKG}/pkg/mqtrigger/kafka -RUN go build -a -o /go/bin/main - -FROM alpine:3.12 as base -RUN apk add --update ca-certificates -COPY --from=builder /go/bin/main / - -ENTRYPOINT ["/main"] \ No newline at end of file diff --git a/pkg/mqtrigger/kafka/main.go b/pkg/mqtrigger/kafka/main.go deleted file mode 100644 index fe1b5dcd..00000000 --- a/pkg/mqtrigger/kafka/main.go +++ /dev/null @@ -1,453 +0,0 @@ -package main - -import ( - "context" - "crypto/sha256" - "crypto/sha512" - "crypto/tls" - "crypto/x509" - "fmt" - "hash" - "io/ioutil" - "log" - "net/http" - "os" - "os/signal" - "strings" - "sync" - "syscall" - "time" - - "github.com/Shopify/sarama" - "github.com/pkg/errors" - "github.com/xdg/scram" - "go.uber.org/zap" - - "github.com/fission/fission/pkg/mqtrigger/util" -) - -type kafkaMetadata struct { - bootstrapServers []string - consumerGroup string - - // auth - authMode kafkaAuthMode - username string - password string - - // ssl - cert string - key string - ca string -} - -type kafkaAuthMode string - -const ( - kafkaAuthModeNone kafkaAuthMode = "none" - kafkaAuthModeSaslPlaintext kafkaAuthMode = "sasl_plaintext" - kafkaAuthModeSaslScramSha256 kafkaAuthMode = "sasl_scram_sha256" - kafkaAuthModeSaslScramSha512 kafkaAuthMode = "sasl_scram_sha512" - kafkaAuthModeSaslSSL kafkaAuthMode = "sasl_ssl" - kafkaAuthModeSaslSSLPlain kafkaAuthMode = "sasl_ssl_plain" -) - -var SHA256 scram.HashGeneratorFcn = func() hash.Hash { return sha256.New() } -var SHA512 scram.HashGeneratorFcn = func() hash.Hash { return sha512.New() } - -type XDGSCRAMClient struct { - *scram.Client - *scram.ClientConversation - scram.HashGeneratorFcn -} - -func (x *XDGSCRAMClient) Begin(userName, password, authzID string) (err error) { - x.Client, err = x.HashGeneratorFcn.NewClient(userName, password, authzID) - if err != nil { - return err - } - x.ClientConversation = x.Client.NewConversation() - return nil -} - -func (x *XDGSCRAMClient) Step(challenge string) (response string, err error) { - response, err = x.ClientConversation.Step(challenge) - return -} - -func (x *XDGSCRAMClient) Done() bool { - return x.ClientConversation.Done() -} - -func parseKafkaMetadata(logger *zap.Logger) (kafkaMetadata, error) { - meta := kafkaMetadata{} - - // brokerList marked as deprecated, bootstrapServers is the new one to use - if os.Getenv("BROKER_LIST") != "" && os.Getenv("BOOTSTRAP_SERVERS") != "" { - return meta, errors.New("cannot specify both bootstrapServers and brokerList (deprecated)") - } - if os.Getenv("BROKER_LIST") == "" && os.Getenv("BOOTSTRAP_SERVERS") == "" { - return meta, errors.New("no bootstrapServers or brokerList (deprecated) given") - } - if os.Getenv("BOOTSTRAP_SERVERS") != "" { - meta.bootstrapServers = strings.Split(os.Getenv("BOOTSTRAP_SERVERS"), ",") - } - if os.Getenv("BROKER_LIST") != "" { - logger.Info("WARNING: usage of brokerList is deprecated. use bootstrapServers instead.") - meta.bootstrapServers = strings.Split(os.Getenv("BROKER_LIST"), ",") - } - if os.Getenv("CONSUMER_GROUP") == "" { - return meta, errors.New("No consumerGroup given") - } - meta.consumerGroup = os.Getenv("CONSUMER_GROUP") - - meta.authMode = kafkaAuthModeNone - mode := kafkaAuthMode(strings.TrimSpace((os.Getenv("AUTH_MODE")))) - if mode == "" { - mode = kafkaAuthModeNone - } - - if mode != kafkaAuthModeNone && mode != kafkaAuthModeSaslPlaintext && mode != kafkaAuthModeSaslSSL && mode != kafkaAuthModeSaslSSLPlain && mode != kafkaAuthModeSaslScramSha256 && mode != kafkaAuthModeSaslScramSha512 { - return meta, fmt.Errorf("err auth mode %s given", mode) - } - - meta.authMode = mode - - if meta.authMode != kafkaAuthModeNone && meta.authMode != kafkaAuthModeSaslSSL { - if os.Getenv("USERNAME") == "" { - return meta, errors.New("no username given") - } - meta.username = strings.TrimSpace(os.Getenv("USERNAME")) - - if os.Getenv("PASSWORD") == "" { - return meta, errors.New("no password given") - } - meta.password = strings.TrimSpace(os.Getenv("PASSWORD")) - } - - if meta.authMode == kafkaAuthModeSaslSSL { - if os.Getenv("CA") == "" { - return meta, errors.New("no ca given") - } - meta.ca = os.Getenv("CA") - - if os.Getenv("CERT") == "" { - return meta, errors.New("no cert given") - } - meta.cert = os.Getenv("CERT") - - if os.Getenv("KEY") == "" { - return meta, errors.New("no key given") - } - meta.key = os.Getenv("KEY") - } - - return meta, nil -} - -func getConfig(metadata kafkaMetadata) (*sarama.Config, error) { - config := sarama.NewConfig() - config.Version = sarama.V1_0_0_0 - - if ok := metadata.authMode == kafkaAuthModeSaslPlaintext || metadata.authMode == kafkaAuthModeSaslSSLPlain || metadata.authMode == kafkaAuthModeSaslScramSha256 || metadata.authMode == kafkaAuthModeSaslScramSha512; ok { - config.Net.SASL.Enable = true - config.Net.SASL.User = metadata.username - config.Net.SASL.Password = metadata.password - } - - if metadata.authMode == kafkaAuthModeSaslSSLPlain { - config.Net.SASL.Mechanism = sarama.SASLMechanism(sarama.SASLTypePlaintext) - - tlsConfig := &tls.Config{ - InsecureSkipVerify: true, - ClientAuth: 0, - } - - config.Net.TLS.Enable = true - config.Net.TLS.Config = tlsConfig - config.Net.DialTimeout = 10 * time.Second - } - - if metadata.authMode == kafkaAuthModeSaslSSL { - cert, err := tls.X509KeyPair([]byte(metadata.cert), []byte(metadata.key)) - if err != nil { - return nil, fmt.Errorf("error parse X509KeyPair: %s", err) - } - - caCertPool := x509.NewCertPool() - caCertPool.AppendCertsFromPEM([]byte(metadata.ca)) - - tlsConfig := &tls.Config{ - Certificates: []tls.Certificate{cert}, - RootCAs: caCertPool, - } - - config.Net.TLS.Enable = true - config.Net.TLS.Config = tlsConfig - } - - if metadata.authMode == kafkaAuthModeSaslScramSha256 { - config.Net.SASL.SCRAMClientGeneratorFunc = func() sarama.SCRAMClient { return &XDGSCRAMClient{HashGeneratorFcn: SHA256} } - config.Net.SASL.Mechanism = sarama.SASLMechanism(sarama.SASLTypeSCRAMSHA256) - } - - if metadata.authMode == kafkaAuthModeSaslScramSha512 { - config.Net.SASL.SCRAMClientGeneratorFunc = func() sarama.SCRAMClient { return &XDGSCRAMClient{HashGeneratorFcn: SHA512} } - config.Net.SASL.Mechanism = sarama.SASLMechanism(sarama.SASLTypeSCRAMSHA512) - } - - if metadata.authMode == kafkaAuthModeSaslPlaintext { - config.Net.SASL.Mechanism = sarama.SASLTypePlaintext - config.Net.TLS.Enable = true - } - return config, nil -} - -// Connector represents a Sarama consumer group consumer -type Connector struct { - ready chan bool - logger *zap.Logger - producer sarama.SyncProducer - fissionTriggerFields util.FissionMetadata -} - -// Setup is run at the beginning of a new session, before ConsumeClaim -func (connector *Connector) Setup(sarama.ConsumerGroupSession) error { - close(connector.ready) - return nil -} - -// Cleanup is run at the end of a session, once all ConsumeClaim goroutines have exited -func (connector *Connector) Cleanup(sarama.ConsumerGroupSession) error { - return nil -} - -// ConsumeClaim must start a consumer loop of ConsumerGroupClaim's Messages() -func (connector *Connector) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error { - for message := range claim.Messages() { - connector.logger.Info(fmt.Sprintf("Message claimed: value = %s, timestamp = %v, topic = %s", string(message.Value), message.Timestamp, message.Topic)) - success := handleFissionFunction(message, connector.fissionTriggerFields, connector.producer, connector.logger) - if success { - session.MarkMessage(message, "") - } - } - return nil -} - -func getProducer(metadata kafkaMetadata) (sarama.SyncProducer, error) { - config, err := getConfig(metadata) - if err != nil { - return nil, err - } - - config.Producer.RequiredAcks = sarama.WaitForAll - config.Producer.Retry.Max = 10 - config.Producer.Return.Successes = true - producer, err := sarama.NewSyncProducer(metadata.bootstrapServers, config) - if err != nil { - return nil, err - } - return producer, nil -} - -func handleFissionFunction(msg *sarama.ConsumerMessage, triggerFields util.FissionMetadata, producer sarama.SyncProducer, logger *zap.Logger) bool { - var value string = string(msg.Value[:]) - // Generate the Headers - fissionHeaders := map[string]string{ - "X-Fission-MQTrigger-Topic": triggerFields.Topic, - "X-Fission-MQTrigger-RespTopic": triggerFields.ResponseTopic, - "X-Fission-MQTrigger-ErrorTopic": triggerFields.ErrorTopic, - "Content-Type": triggerFields.ContentType, - } - - // Create request - req, err := http.NewRequest("POST", triggerFields.FunctionURL, strings.NewReader(value)) - if err != nil { - logger.Error("failed to create HTTP request to invoke function", - zap.Error(err), - zap.String("function_url", triggerFields.FunctionURL)) - return false - } - - // Set the headers came from Kafka record - // Using Header.Add() as msg.Headers may have keys with more than one value - for _, h := range msg.Headers { - req.Header.Add(string(h.Key), string(h.Value)) - } - - for k, v := range fissionHeaders { - req.Header.Set(k, v) - } - - // Make the request - var resp *http.Response - for attempt := 0; attempt <= triggerFields.MaxRetries; attempt++ { - // Make the request - resp, err = http.DefaultClient.Do(req) - if err != nil { - logger.Error("sending function invocation request failed", - zap.Error(err), - zap.String("function_url", triggerFields.FunctionURL), - zap.String("trigger", triggerFields.TriggerName)) - continue - } - if resp == nil { - continue - } - if err == nil && resp.StatusCode == http.StatusOK { - // Success, quit retrying - break - } - } - - if resp == nil { - logger.Warn("every function invocation retry failed; final retry gave empty response", - zap.String("function_url", triggerFields.FunctionURL), - zap.String("trigger", triggerFields.TriggerName)) - return false - } - defer resp.Body.Close() - body, err := ioutil.ReadAll(resp.Body) - - logger.Debug("got response from function invocation", - zap.String("function_url", triggerFields.FunctionURL), - zap.String("trigger", triggerFields.TriggerName), - zap.String("body", string(body))) - - if err != nil { - errorHandler(logger, triggerFields, producer, - errors.Wrapf(err, "request body error: %v", string(body))) - return false - } - if resp.StatusCode != 200 { - errorHandler(logger, triggerFields, producer, - fmt.Errorf("request returned failure: %v", resp.StatusCode)) - return false - } - - if len(triggerFields.ResponseTopic) > 0 { - // Generate Kafka record headers - var kafkaRecordHeaders []sarama.RecordHeader - - for k, v := range resp.Header { - // One key may have multiple values - for _, v := range v { - kafkaRecordHeaders = append(kafkaRecordHeaders, sarama.RecordHeader{Key: []byte(k), Value: []byte(v)}) - } - } - - _, _, err = producer.SendMessage(&sarama.ProducerMessage{ - Topic: triggerFields.ResponseTopic, - Value: sarama.StringEncoder(body), - Headers: kafkaRecordHeaders, - }) - if err != nil { - logger.Warn("failed to publish response body from function invocation to topic", - zap.Error(err), - zap.String("topic", triggerFields.Topic), - zap.String("function_url", triggerFields.FunctionURL)) - return false - } - } - - return true -} - -func errorHandler(logger *zap.Logger, triggerFields util.FissionMetadata, producer sarama.SyncProducer, err error) { - if len(triggerFields.ErrorTopic) > 0 { - _, _, e := producer.SendMessage(&sarama.ProducerMessage{ - Topic: triggerFields.ErrorTopic, - Value: sarama.StringEncoder(err.Error()), - }) - if e != nil { - logger.Error("failed to publish message to error topic", - zap.Error(e), - zap.String("trigger", triggerFields.TriggerName), - zap.String("message", err.Error()), - zap.String("topic", triggerFields.Topic)) - } - } else { - logger.Error("message received to publish to error topic, but no error topic was set", - zap.String("message", err.Error()), zap.String("trigger", triggerFields.TriggerName), zap.String("function_url", triggerFields.FunctionURL)) - } -} - -func main() { - logger, err := zap.NewProduction() - if err != nil { - log.Fatalf("can't initialize zap logger: %v", err) - } - defer logger.Sync() - - metadata, err := parseKafkaMetadata(logger) - if err != nil { - logger.Error("Failed to fetch kafka metadata", zap.Error(err)) - return - } - - triggerFields, err := util.ParseFissionMetadata() - if err != nil { - logger.Error("Failed to parse fission trigger fields", zap.Error(err)) - return - } - - config, err := getConfig(metadata) - if err != nil { - logger.Error("Failed to create kafka config", zap.Error(err)) - return - } - - producer, err := getProducer(metadata) - if err != nil { - logger.Error("Failed to create kafka producer", zap.Error(err)) - return - } - defer producer.Close() - - connector := Connector{ - ready: make(chan bool), - logger: logger, - producer: producer, - fissionTriggerFields: triggerFields, - } - - ctx, cancel := context.WithCancel(context.Background()) - client, err := sarama.NewConsumerGroup(metadata.bootstrapServers, metadata.consumerGroup, config) - if err != nil { - logger.Error("Error creating consumer group client", zap.Error(err)) - return - } - - wg := &sync.WaitGroup{} - wg.Add(1) - - go func() { - defer wg.Done() - for { - if err := client.Consume(ctx, []string{triggerFields.Topic}, &connector); err != nil { - logger.Error("Error from consumer", zap.Error(err)) - } - // check if context was cancelled, signaling that the consumer should stop - if ctx.Err() != nil { - return - } - connector.ready = make(chan bool) - } - }() - - <-connector.ready // Await till the consumer has been set up - logger.Info("Sarama consumer up and running!...") - sigterm := make(chan os.Signal, 1) - signal.Notify(sigterm, syscall.SIGINT, syscall.SIGTERM) - select { - case <-ctx.Done(): - logger.Info("terminating: context cancelled") - case <-sigterm: - logger.Info("terminating: via signal") - } - cancel() - wg.Wait() - if err = client.Close(); err != nil { - logger.Error("Error closing client", zap.Error(err)) - } -} diff --git a/pkg/mqtrigger/scalermanager.go b/pkg/mqtrigger/scalermanager.go index 3c104c35..2de1bbb6 100644 --- a/pkg/mqtrigger/scalermanager.go +++ b/pkg/mqtrigger/scalermanager.go @@ -54,7 +54,7 @@ func getAuthTriggerClient(namespace string) (dynamic.ResourceInterface, error) { if err != nil { return nil, err } - return dynamicClient.Resource(scaledObjectGVR).Namespace(namespace), nil + return dynamicClient.Resource(authTriggerGVR).Namespace(namespace), nil } // StartScalerManager watches for changes in MessageQueueTrigger and, @@ -166,7 +166,7 @@ func getEnvVarlist(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient ku Value: mqt.Spec.Topic, }, { - Name: "FUNCTION_URL", + Name: "HTTP_ENDPOINT", Value: url, }, { @@ -178,7 +178,7 @@ func getEnvVarlist(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient ku Value: mqt.Spec.ResponseTopic, }, { - Name: "TRIGGER_NAME", + Name: "SOURCE_NAME", Value: mqt.ObjectMeta.Name, }, { diff --git a/pkg/mqtrigger/scalermanager_test.go b/pkg/mqtrigger/scalermanager_test.go index 592c070d..9ee0a7f7 100644 --- a/pkg/mqtrigger/scalermanager_test.go +++ b/pkg/mqtrigger/scalermanager_test.go @@ -111,7 +111,7 @@ func Test_getEnvVarlist(t *testing.T) { Value: mqt.Spec.Topic, }, { - Name: "FUNCTION_URL", + Name: "HTTP_ENDPOINT", Value: "http://router.fission/fission-function/fission-function/test", }, { @@ -123,7 +123,7 @@ func Test_getEnvVarlist(t *testing.T) { Value: "response-topic", }, { - Name: "TRIGGER_NAME", + Name: "SOURCE_NAME", Value: "Test", }, { diff --git a/pkg/mqtrigger/util/util.go b/pkg/mqtrigger/util/util.go deleted file mode 100644 index e1665f2b..00000000 --- a/pkg/mqtrigger/util/util.go +++ /dev/null @@ -1,43 +0,0 @@ -package util - -import ( - "fmt" - "os" - "strconv" - "strings" -) - -// FissionMetadata contains common fission side fields -type FissionMetadata struct { - // fission - Topic string - ResponseTopic string - ErrorTopic string - FunctionURL string - MaxRetries int - ContentType string - TriggerName string -} - -// ParseFissionMetadata parses fission side common fields and returns as fissionMetadata or returns error -func ParseFissionMetadata() (FissionMetadata, error) { - for _, envVars := range []string{"TOPIC", "FUNCTION_URL", "MAX_RETRIES", "CONTENT_TYPE", "TRIGGER_NAME"} { - if os.Getenv(envVars) == "" { - return FissionMetadata{}, fmt.Errorf("environment variable not found: %v", envVars) - } - } - meta := FissionMetadata{ - Topic: os.Getenv("TOPIC"), - ResponseTopic: os.Getenv("RESPONSE_TOPIC"), - ErrorTopic: os.Getenv("ERROR_TOPIC"), - FunctionURL: os.Getenv("FUNCTION_URL"), - ContentType: os.Getenv("CONTENT_TYPE"), - TriggerName: os.Getenv("TRIGGER_NAME"), - } - val, err := strconv.ParseInt(strings.TrimSpace(os.Getenv("MAX_RETRIES")), 0, 64) - if err != nil { - return FissionMetadata{}, fmt.Errorf("failed to parse value from MAX_RETRIES environment variable %v", err) - } - meta.MaxRetries = int(val) - return meta, nil -} diff --git a/pkg/mqtrigger/validator/validator.go b/pkg/mqtrigger/validator/validator.go index fddcbe3c..f5d2f0cb 100644 --- a/pkg/mqtrigger/validator/validator.go +++ b/pkg/mqtrigger/validator/validator.go @@ -45,7 +45,10 @@ func Register(mqType string, validator TopicValidator) { topicValidators[mqType] = validator } -func IsValidTopic(mqType string, topic string) bool { +func IsValidTopic(mqType, topic, mqtKind string) bool { + if mqtKind == "keda" { + return true + } validator, registered := topicValidators[mqType] if !registered { return false @@ -53,7 +56,10 @@ func IsValidTopic(mqType string, topic string) bool { return validator(topic) } -func IsValidMessageQueue(mqType string) bool { +func IsValidMessageQueue(mqType, mqtKind string) bool { + if mqtKind == "keda" { + return true + } _, registered := topicValidators[mqType] return registered }