diff --git a/cmd/fission-bundle/mqtrigger/mqtrigger.go b/cmd/fission-bundle/mqtrigger/mqtrigger.go index 57595755..31b2589e 100644 --- a/cmd/fission-bundle/mqtrigger/mqtrigger.go +++ b/cmd/fission-bundle/mqtrigger/mqtrigger.go @@ -29,10 +29,11 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/mqtrigger" + "github.com/fission/fission/pkg/mqtrigger/factory" "github.com/fission/fission/pkg/mqtrigger/messageQueue" - "github.com/fission/fission/pkg/mqtrigger/messageQueue/azurequeuestorage" - "github.com/fission/fission/pkg/mqtrigger/messageQueue/kafka" - "github.com/fission/fission/pkg/mqtrigger/messageQueue/nats" + _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/azurequeuestorage" + _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/kafka" + _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/nats" ) func Start(logger *zap.Logger, routerUrl string) error { @@ -47,8 +48,7 @@ func Start(logger *zap.Logger, routerUrl string) error { return errors.Wrap(err, "error waiting for CRDs") } - // Message queue type: nats is the only supported one for now - mqType := os.Getenv("MESSAGE_QUEUE_TYPE") + mqType := (fv1.MessageQueueType)(os.Getenv("MESSAGE_QUEUE_TYPE")) mqUrl := os.Getenv("MESSAGE_QUEUE_URL") secretsPath := strings.TrimSpace(os.Getenv("MESSAGE_QUEUE_SECRETS")) @@ -62,38 +62,25 @@ func Start(logger *zap.Logger, routerUrl string) error { } } - mq, err := newMessageQueue( + mq, err := factory.Create( logger, - routerUrl, + mqType, messageQueue.Config{ - MQType: mqType, + MQType: (string)(mqType), Url: mqUrl, Secrets: secrets, }, + routerUrl, ) if err != nil { logger.Fatal("failed to connect to remote message queue server", zap.Error(err)) } - mqtrigger.MakeMessageQueueTriggerManager(logger, fissionClient, mq).Run() + mqtrigger.MakeMessageQueueTriggerManager(logger, fissionClient, mqType, mq).Run() return nil } -func newMessageQueue(logger *zap.Logger, routerURL string, mqCfg messageQueue.Config) (messageQueue messageQueue.MessageQueue, err error) { - switch mqCfg.MQType { - case fv1.MessageQueueTypeNats: - messageQueue, err = nats.New(logger, routerURL, mqCfg) - case fv1.MessageQueueTypeASQ: - messageQueue, err = azurequeuestorage.New(logger, routerURL, mqCfg) - case fv1.MessageQueueTypeKafka: - messageQueue, err = kafka.New(logger, routerURL, mqCfg) - default: - err = errors.Errorf("no supported message queue type found for %q", mqCfg.MQType) - } - return messageQueue, err -} - func readSecrets(logger *zap.Logger, secretsPath string) (map[string][]byte, error) { // return if no secrets exist if _, err := os.Stat(secretsPath); os.IsNotExist(err) { diff --git a/cmd/fission-cli/app/app.go b/cmd/fission-cli/app/app.go index 7b42b952..1f304606 100644 --- a/cmd/fission-cli/app/app.go +++ b/cmd/fission-cli/app/app.go @@ -24,6 +24,9 @@ import ( "github.com/fission/fission/pkg/fission-cli/flag" flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" "github.com/fission/fission/pkg/fission-cli/util" + _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/azurequeuestorage" + _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/kafka" + _ "github.com/fission/fission/pkg/mqtrigger/messageQueue/nats" ) const ( diff --git a/pkg/apis/core/v1/validation.go b/pkg/apis/core/v1/validation.go index cd5e286d..98a76ac0 100644 --- a/pkg/apis/core/v1/validation.go +++ b/pkg/apis/core/v1/validation.go @@ -23,10 +23,11 @@ import ( "strings" "github.com/hashicorp/go-multierror" - nsUtil "github.com/nats-io/nats-streaming-server/util" "github.com/robfig/cron" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/validation" + + "github.com/fission/fission/pkg/mqtrigger/validator" ) const ( @@ -37,12 +38,6 @@ const ( totalAnnotationSizeLimitB int = 256 * (1 << 10) // 256 kB ) -var ( - validAzureQueueName = regexp.MustCompile(`^[a-z0-9][a-z0-9\\-]*[a-z0-9]$`) - // Need to use raw string to support escape sequence for - & . chars - validKafkaTopicName = regexp.MustCompile(`^[a-zA-Z0-9][a-zA-Z0-9\-\._]*[a-zA-Z0-9]$`) -) - type ( ValidationErrorType int @@ -163,35 +158,6 @@ func ValidateKubeReference(refName string, name string, namespace string) error return result.ErrorOrNil() } -func IsTopicValid(mqType MessageQueueType, topic string) bool { - switch mqType { - case MessageQueueTypeNats: - return nsUtil.IsChannelNameValid(topic, false) - case MessageQueueTypeASQ: - return len(topic) >= 3 && len(topic) <= 63 && validAzureQueueName.MatchString(topic) - case MessageQueueTypeKafka: - return IsValidKafkaTopic(topic) - } - return false -} - -// The validation is based on Kafka's internal implementation: https://github.com/apache/kafka/blob/trunk/clients/src/main/java/org/apache/kafka/common/internals/Topic.java -func IsValidKafkaTopic(topic string) bool { - if len(topic) == 0 { - return false - } - if topic == "." || topic == ".." { - return false - } - if len(topic) > 249 { - return false - } - if !validKafkaTopicName.MatchString(topic) { - return false - } - return true -} - func IsValidCronSpec(spec string) error { _, err := cron.Parse(spec) return err @@ -528,18 +494,16 @@ func (spec MessageQueueTriggerSpec) Validate() error { result = multierror.Append(result, spec.FunctionReference.Validate()) - switch spec.MessageQueueType { - case MessageQueueTypeNats, MessageQueueTypeASQ, MessageQueueTypeKafka: // no op - default: + if !validator.IsValidMessageQueue((string)(spec.MessageQueueType)) { 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) { + result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "MessageQueueTriggerSpec.Topic", spec.Topic, "not a valid topic")) + } - if !IsTopicValid(spec.MessageQueueType, spec.Topic) { - result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "MessageQueueTriggerSpec.Topic", spec.Topic, "not a valid topic")) - } - - if len(spec.ResponseTopic) > 0 && !IsTopicValid(spec.MessageQueueType, spec.ResponseTopic) { - result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "MessageQueueTriggerSpec.ResponseTopic", spec.ResponseTopic, "not a valid topic")) + if len(spec.ResponseTopic) > 0 && !validator.IsValidTopic((string)(spec.MessageQueueType), spec.ResponseTopic) { + result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "MessageQueueTriggerSpec.ResponseTopic", spec.ResponseTopic, "not a valid topic")) + } } return result.ErrorOrNil() diff --git a/pkg/fission-cli/cmd/mqtrigger/create.go b/pkg/fission-cli/cmd/mqtrigger/create.go index ef5a64f3..71e90ad3 100644 --- a/pkg/fission-cli/cmd/mqtrigger/create.go +++ b/pkg/fission-cli/cmd/mqtrigger/create.go @@ -30,6 +30,7 @@ import ( "github.com/fission/fission/pkg/fission-cli/console" flagkey "github.com/fission/fission/pkg/fission-cli/flag/key" "github.com/fission/fission/pkg/fission-cli/util" + "github.com/fission/fission/pkg/mqtrigger/validator" ) type CreateSubCommand struct { @@ -58,18 +59,9 @@ func (opts *CreateSubCommand) complete(input cli.Input) error { fnName := input.String(flagkey.MqtFnName) fnNamespace := input.String(flagkey.NamespaceFunction) - var mqType fv1.MessageQueueType - switch input.String(flagkey.MqtMQType) { - case "": - mqType = fv1.MessageQueueTypeNats - case fv1.MessageQueueTypeNats: - mqType = fv1.MessageQueueTypeNats - case fv1.MessageQueueTypeASQ: - mqType = fv1.MessageQueueTypeASQ - case fv1.MessageQueueTypeKafka: - mqType = fv1.MessageQueueTypeKafka - default: - return errors.New("Unknown message queue type, currently only \"nats-streaming, azure-storage-queue, kafka \" is supported") + mqType := (fv1.MessageQueueType)(input.String(flagkey.MqtMQType)) + if !validator.IsValidMessageQueue((string)(mqType)) { + return errors.New("Unsupported message queue type") } topic := input.String(flagkey.MqtTopic) @@ -172,7 +164,7 @@ func (opts *CreateSubCommand) run(input cli.Input) error { func checkMQTopicAvailability(mqType fv1.MessageQueueType, topics ...string) error { for _, t := range topics { - if len(t) > 0 && !fv1.IsTopicValid(mqType, t) { + if len(t) > 0 && !validator.IsValidTopic((string)(mqType), t) { return errors.Errorf("invalid topic for %s: %s", mqType, t) } } diff --git a/pkg/fission-cli/flag/flag.go b/pkg/fission-cli/flag/flag.go index 73734b31..3d883b34 100644 --- a/pkg/fission-cli/flag/flag.go +++ b/pkg/fission-cli/flag/flag.go @@ -127,7 +127,7 @@ var ( MqtName = Flag{Type: String, Name: flagkey.MqtName, Usage: "Message queue trigger name"} MqtFnName = Flag{Type: String, Name: flagkey.MqtFnName, Usage: "Function name"} - MqtMQType = Flag{Type: String, Name: flagkey.MqtMQType, Usage: "Message queue type, e.g. nats-streaming, azure-storage-queue", DefaultValue: "nats-streaming"} + MqtMQType = Flag{Type: String, Name: flagkey.MqtMQType, Usage: "Message queue type, e.g. nats-streaming, azure-storage-queue, kafka", DefaultValue: "nats-streaming"} MqtTopic = Flag{Type: String, Name: flagkey.MqtTopic, Usage: "Message queue Topic the trigger listens on"} MqtRespTopic = Flag{Type: String, Name: flagkey.MqtRespTopic, Usage: "Topic that the function response is sent on (response discarded if unspecified)"} MqtErrorTopic = Flag{Type: String, Name: flagkey.MqtErrorTopic, Usage: "Topic that the function error messages are sent to (errors discarded if unspecified"} diff --git a/pkg/mqtrigger/factory/factory.go b/pkg/mqtrigger/factory/factory.go new file mode 100644 index 00000000..1bc79b1d --- /dev/null +++ b/pkg/mqtrigger/factory/factory.go @@ -0,0 +1,62 @@ +/* +Copyright 2020 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package factory + +import ( + "sync" + + "github.com/pkg/errors" + "go.uber.org/zap" + + fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/mqtrigger/messageQueue" +) + +var ( + messageQueueFactories = make(map[fv1.MessageQueueType]MessageQueueFactory) + lock = sync.Mutex{} +) + +type ( + MessageQueueFactory interface { + Create(logger *zap.Logger, config messageQueue.Config, routerURL string) (messageQueue.MessageQueue, error) + } +) + +func Register(mqType fv1.MessageQueueType, factory MessageQueueFactory) { + lock.Lock() + defer lock.Unlock() + + if factory == nil { + panic("Nil message queue factory") + } + + _, registered := messageQueueFactories[mqType] + if registered { + panic("Message queue factory already register") + } + + messageQueueFactories[mqType] = factory +} + +func Create(logger *zap.Logger, mqType fv1.MessageQueueType, mqConfig messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { + factory, registered := messageQueueFactories[mqType] + if !registered { + return nil, errors.Errorf("no supported message queue type found for %q", mqType) + } + return factory.Create(logger, mqConfig, routerUrl) +} diff --git a/pkg/mqtrigger/messageQueue/azurequeuestorage/asq.go b/pkg/mqtrigger/messageQueue/azurequeuestorage/asq.go index 8833b10b..78426b39 100644 --- a/pkg/mqtrigger/messageQueue/azurequeuestorage/asq.go +++ b/pkg/mqtrigger/messageQueue/azurequeuestorage/asq.go @@ -30,14 +30,21 @@ import ( "time" "github.com/Azure/azure-sdk-for-go/storage" - "github.com/fission/fission/pkg/utils" "github.com/pkg/errors" "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/mqtrigger/factory" "github.com/fission/fission/pkg/mqtrigger/messageQueue" + "github.com/fission/fission/pkg/mqtrigger/validator" + "github.com/fission/fission/pkg/utils" ) +func init() { + factory.Register(fv1.MessageQueueTypeASQ, &Factory{}) + validator.Register(fv1.MessageQueueTypeASQ, IsTopicValid) +} + // TODO: some of these constants should probably be environment variables const ( // AzureQueuePollingInterval is the polling interval (default is 1 minute). @@ -105,6 +112,12 @@ type AzureHTTPClient interface { Do(req *http.Request) (*http.Response, error) } +type Factory struct{} + +func (factory *Factory) Create(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { + return New(logger, mqCfg, routerUrl) +} + type azureQueueService struct { service storage.QueueServiceClient } @@ -179,7 +192,7 @@ func newAzureQueueService(client storage.Client) AzureQueueService { } } -func New(logger *zap.Logger, routerURL string, config messageQueue.Config) (messageQueue.MessageQueue, error) { +func New(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { account := os.Getenv("AZURE_STORAGE_ACCOUNT_NAME") if len(account) == 0 { return nil, errors.New("Required environment variable 'AZURE_STORAGE_ACCOUNT_NAME' is not set") @@ -198,7 +211,7 @@ func New(logger *zap.Logger, routerURL string, config messageQueue.Config) (mess } return &AzureStorageConnection{ logger: logger.Named("azue_storage"), - routerURL: routerURL, + routerURL: routerUrl, service: newAzureQueueService(client), httpClient: &http.Client{ Timeout: AzureFunctionInvocationTimeout, diff --git a/pkg/mqtrigger/messageQueue/azurequeuestorage/asq_test.go b/pkg/mqtrigger/messageQueue/azurequeuestorage/asq_test.go index fa4af81c..6bd75bac 100644 --- a/pkg/mqtrigger/messageQueue/azurequeuestorage/asq_test.go +++ b/pkg/mqtrigger/messageQueue/azurequeuestorage/asq_test.go @@ -113,10 +113,10 @@ func TestNewStorageConnectionMissingAccountName(t *testing.T) { logger, err := zap.NewDevelopment() panicIf(err) - connection, err := New(logger, DummyRouterURL, messageQueue.Config{ + connection, err := New(logger, messageQueue.Config{ MQType: fv1.MessageQueueTypeASQ, Url: "", - }) + }, DummyRouterURL) require.Nil(t, connection) require.Error(t, err, "Required environment variable 'AZURE_STORAGE_ACCOUNT_NAME' is not set") } @@ -126,10 +126,10 @@ func TestNewStorageConnectionMissingAccessKey(t *testing.T) { panicIf(err) _ = os.Setenv("AZURE_STORAGE_ACCOUNT_NAME", "accountname") - connection, err := New(logger, DummyRouterURL, messageQueue.Config{ + connection, err := New(logger, messageQueue.Config{ MQType: fv1.MessageQueueTypeASQ, Url: "", - }) + }, DummyRouterURL) _ = os.Unsetenv("AZURE_STORAGE_ACCOUNT_NAME") require.Nil(t, connection) require.Error(t, err, "Required environment variable 'AZURE_STORAGE_ACCOUNT_KEY' is not set") @@ -141,10 +141,10 @@ func TestNewStorageConnection(t *testing.T) { _ = os.Setenv("AZURE_STORAGE_ACCOUNT_NAME", "accountname") _ = os.Setenv("AZURE_STORAGE_ACCOUNT_KEY", "bm90IGEga2V5") - connection, err := New(logger, DummyRouterURL, messageQueue.Config{ + connection, err := New(logger, messageQueue.Config{ MQType: "azure-storage-queue", Url: "", - }) + }, DummyRouterURL) _ = os.Unsetenv("AZURE_STORAGE_ACCOUNT_NAME") _ = os.Unsetenv("AZURE_STORAGE_ACCOUNT_KEY") require.NoError(t, err) diff --git a/pkg/mqtrigger/messageQueue/kafka/kafka.go b/pkg/mqtrigger/messageQueue/kafka/kafka.go index d54d5ba7..8541a821 100644 --- a/pkg/mqtrigger/messageQueue/kafka/kafka.go +++ b/pkg/mqtrigger/messageQueue/kafka/kafka.go @@ -33,10 +33,17 @@ import ( "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/mqtrigger/factory" "github.com/fission/fission/pkg/mqtrigger/messageQueue" + "github.com/fission/fission/pkg/mqtrigger/validator" "github.com/fission/fission/pkg/utils" ) +func init() { + factory.Register(fv1.MessageQueueTypeKafka, &Factory{}) + validator.Register(fv1.MessageQueueTypeKafka, IsTopicValid) +} + var ( // Need to use raw string to support escape sequence for - & . chars validKafkaTopicName = regexp.MustCompile(`^[a-zA-Z0-9][a-zA-Z0-9\-\._]*[a-zA-Z0-9]$`) @@ -51,9 +58,15 @@ type ( authKeys map[string][]byte tls bool } + + Factory struct{} ) -func New(logger *zap.Logger, routerUrl string, mqCfg messageQueue.Config) (messageQueue.MessageQueue, error) { +func (factory *Factory) Create(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { + return New(logger, mqCfg, routerUrl) +} + +func New(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { if len(routerUrl) == 0 || len(mqCfg.Url) == 0 { return nil, errors.New("the router URL or MQ URL is empty") } @@ -340,7 +353,8 @@ func errorHandler(logger *zap.Logger, trigger *fv1.MessageQueueTrigger, producer } } -// The validation is based on Kafka's internal implementation: https://github.com/apache/kafka/blob/cde6d18983b5d58199f8857d8d61d7efcbe6e54a/clients/src/main/java/org/apache/kafka/common/internals/Topic.java#L36-L47 +// The validation is based on Kafka's internal implementation: +// https://github.com/apache/kafka/blob/cde6d18983b5d58199f8857d8d61d7efcbe6e54a/clients/src/main/java/org/apache/kafka/common/internals/Topic.java#L36-L47 func IsTopicValid(topic string) bool { if len(topic) == 0 { return false diff --git a/pkg/mqtrigger/messageQueue/nats/nats.go b/pkg/mqtrigger/messageQueue/nats/nats.go index 713b707a..b4cacab9 100644 --- a/pkg/mqtrigger/messageQueue/nats/nats.go +++ b/pkg/mqtrigger/messageQueue/nats/nats.go @@ -28,10 +28,17 @@ import ( "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/mqtrigger/factory" "github.com/fission/fission/pkg/mqtrigger/messageQueue" + "github.com/fission/fission/pkg/mqtrigger/validator" "github.com/fission/fission/pkg/utils" ) +func init() { + factory.Register(fv1.MessageQueueTypeNats, &Factory{}) + validator.Register(fv1.MessageQueueTypeNats, IsTopicValid) +} + const ( natsClusterID = "fissionMQTrigger" natsProtocol = "nats://" @@ -45,9 +52,15 @@ type ( nsConn ns.Conn routerUrl string } + + Factory struct{} ) -func New(logger *zap.Logger, routerUrl string, mqCfg messageQueue.Config) (messageQueue.MessageQueue, error) { +func (factory *Factory) Create(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { + return New(logger, mqCfg, routerUrl) +} + +func New(logger *zap.Logger, mqCfg messageQueue.Config, routerUrl string) (messageQueue.MessageQueue, error) { conn, err := ns.Connect(natsClusterID, natsClientID, ns.NatsURL(mqCfg.Url), ns.SetConnectionLostHandler(func(conn ns.Conn, reason error) { // TODO: Better way to handle connection lost problem. diff --git a/pkg/mqtrigger/mqtmanager.go b/pkg/mqtrigger/mqtmanager.go index 0ec670f3..0fbb78b0 100644 --- a/pkg/mqtrigger/mqtmanager.go +++ b/pkg/mqtrigger/mqtmanager.go @@ -26,9 +26,6 @@ import ( fv1 "github.com/fission/fission/pkg/apis/core/v1" "github.com/fission/fission/pkg/crd" "github.com/fission/fission/pkg/mqtrigger/messageQueue" - "github.com/fission/fission/pkg/mqtrigger/messageQueue/azurequeuestorage" - "github.com/fission/fission/pkg/mqtrigger/messageQueue/kafka" - "github.com/fission/fission/pkg/mqtrigger/messageQueue/nats" "github.com/fission/fission/pkg/utils" ) @@ -42,11 +39,12 @@ type ( requestType int MessageQueueTriggerManager struct { - logger *zap.Logger - reqChan chan request - triggers map[string]*triggerSubscription - fissionClient *crd.FissionClient - messageQueue messageQueue.MessageQueue + logger *zap.Logger + reqChan chan request + triggers map[string]*triggerSubscription + fissionClient *crd.FissionClient + messageQueueType fv1.MessageQueueType + messageQueue messageQueue.MessageQueue } triggerSubscription struct { @@ -66,13 +64,14 @@ type ( ) func MakeMessageQueueTriggerManager(logger *zap.Logger, - fissionClient *crd.FissionClient, messageQueue messageQueue.MessageQueue) *MessageQueueTriggerManager { + fissionClient *crd.FissionClient, mqType fv1.MessageQueueType, messageQueue messageQueue.MessageQueue) *MessageQueueTriggerManager { mqTriggerMgr := MessageQueueTriggerManager{ - logger: logger.Named("message_queue_trigger_manager"), - reqChan: make(chan request), - triggers: make(map[string]*triggerSubscription), - fissionClient: fissionClient, - messageQueue: messageQueue, + logger: logger.Named("message_queue_trigger_manager"), + reqChan: make(chan request), + triggers: make(map[string]*triggerSubscription), + fissionClient: fissionClient, + messageQueueType: mqType, + messageQueue: messageQueue, } return &mqTriggerMgr } @@ -154,7 +153,9 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() { newTriggerMap := make(map[string]*fv1.MessageQueueTrigger) for index := range newTriggers.Items { newTrigger := &newTriggers.Items[index] - newTriggerMap[crd.CacheKey(&newTrigger.ObjectMeta)] = newTrigger + if newTrigger.Spec.MessageQueueType == mqt.messageQueueType { + newTriggerMap[crd.CacheKey(&newTrigger.ObjectMeta)] = newTrigger + } } // get current set of triggers @@ -205,15 +206,3 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() { time.Sleep(3 * time.Second) } } - -func IsTopicValid(mqType fv1.MessageQueueType, topic string) bool { - switch mqType { - case fv1.MessageQueueTypeNats: - return nats.IsTopicValid(topic) - case fv1.MessageQueueTypeASQ: - return azurequeuestorage.IsTopicValid(topic) - case fv1.MessageQueueTypeKafka: - return kafka.IsTopicValid(topic) - } - return false -} diff --git a/pkg/mqtrigger/validator/validator.go b/pkg/mqtrigger/validator/validator.go new file mode 100644 index 00000000..fddcbe3c --- /dev/null +++ b/pkg/mqtrigger/validator/validator.go @@ -0,0 +1,59 @@ +/* +Copyright 2020 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package validator + +import ( + "sync" +) + +var ( + topicValidators = make(map[string]TopicValidator) + lock = sync.Mutex{} +) + +type ( + TopicValidator func(topic string) bool +) + +func Register(mqType string, validator TopicValidator) { + lock.Lock() + defer lock.Unlock() + + if validator == nil { + panic("Nil message queue topic validator") + } + + _, registered := topicValidators[mqType] + if registered { + panic("Message queue topic validator already register") + } + + topicValidators[mqType] = validator +} + +func IsValidTopic(mqType string, topic string) bool { + validator, registered := topicValidators[mqType] + if !registered { + return false + } + return validator(topic) +} + +func IsValidMessageQueue(mqType string) bool { + _, registered := topicValidators[mqType] + return registered +}