diff --git a/cmd/fission-bundle/main.go b/cmd/fission-bundle/main.go index 1acdf35a..52b3ccda 100644 --- a/cmd/fission-bundle/main.go +++ b/cmd/fission-bundle/main.go @@ -28,13 +28,13 @@ import ( "go.opencensus.io/trace" "go.uber.org/zap" + "github.com/fission/fission/cmd/fission-bundle/mqtrigger" "github.com/fission/fission/pkg/buildermgr" "github.com/fission/fission/pkg/controller" "github.com/fission/fission/pkg/executor" "github.com/fission/fission/pkg/info" "github.com/fission/fission/pkg/kubewatcher" functionLogger "github.com/fission/fission/pkg/logger" - messagequeue "github.com/fission/fission/pkg/mqtrigger" "github.com/fission/fission/pkg/router" "github.com/fission/fission/pkg/storagesvc" "github.com/fission/fission/pkg/timer" @@ -72,7 +72,7 @@ func runTimer(logger *zap.Logger, routerUrl string) { } func runMessageQueueMgr(logger *zap.Logger, routerUrl string) { - err := messagequeue.Start(logger, routerUrl) + err := mqtrigger.Start(logger, routerUrl) if err != nil { logger.Fatal("error starting message queue manager", zap.Error(err)) } diff --git a/pkg/mqtrigger/mqtrigger.go b/cmd/fission-bundle/mqtrigger/mqtrigger.go similarity index 66% rename from pkg/mqtrigger/mqtrigger.go rename to cmd/fission-bundle/mqtrigger/mqtrigger.go index 9012cf15..57595755 100644 --- a/pkg/mqtrigger/mqtrigger.go +++ b/cmd/fission-bundle/mqtrigger/mqtrigger.go @@ -26,8 +26,13 @@ import ( "github.com/pkg/errors" "go.uber.org/zap" + 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/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" ) func Start(logger *zap.Logger, routerUrl string) error { @@ -57,17 +62,39 @@ func Start(logger *zap.Logger, routerUrl string) error { } } - mqCfg := messageQueue.MessageQueueConfig{ - MQType: mqType, - Url: mqUrl, - Secrets: secrets, + mq, err := newMessageQueue( + logger, + routerUrl, + messageQueue.Config{ + MQType: mqType, + Url: mqUrl, + Secrets: secrets, + }, + ) + if err != nil { + logger.Fatal("failed to connect to remote message queue server", zap.Error(err)) } - messageQueue.MakeMessageQueueTriggerManager(logger, fissionClient, routerUrl, mqCfg) + + mqtrigger.MakeMessageQueueTriggerManager(logger, fissionClient, mq).Run() + return nil } -func readSecrets(logger *zap.Logger, secretsPath string) (map[string][]byte, error) { +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) { return nil, err diff --git a/pkg/mqtrigger/messageQueue/asq.go b/pkg/mqtrigger/messageQueue/azurequeuestorage/asq.go similarity index 95% rename from pkg/mqtrigger/messageQueue/asq.go rename to pkg/mqtrigger/messageQueue/azurequeuestorage/asq.go index f6911389..8833b10b 100644 --- a/pkg/mqtrigger/messageQueue/asq.go +++ b/pkg/mqtrigger/messageQueue/azurequeuestorage/asq.go @@ -14,7 +14,7 @@ See the License for the specific language governing permissions and limitations under the License. */ -package messageQueue +package azurequeuestorage import ( "bytes" @@ -23,6 +23,7 @@ import ( "io/ioutil" "net/http" "os" + "regexp" "strconv" "strings" "sync" @@ -34,6 +35,7 @@ import ( "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/mqtrigger/messageQueue" ) // TODO: some of these constants should probably be environment variables @@ -52,6 +54,10 @@ const ( AzureFunctionInvocationTimeout = 10 * time.Minute ) +var ( + validAzureQueueName = regexp.MustCompile(`^[a-z0-9][a-z0-9\\-]*[a-z0-9]$`) +) + // AzureStorageConnection represents an Azure storage connection. type AzureStorageConnection struct { logger *zap.Logger @@ -173,7 +179,7 @@ func newAzureQueueService(client storage.Client) AzureQueueService { } } -func newAzureStorageConnection(logger *zap.Logger, routerURL string, config MessageQueueConfig) (MessageQueue, error) { +func New(logger *zap.Logger, routerURL string, config messageQueue.Config) (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") @@ -200,7 +206,7 @@ func newAzureStorageConnection(logger *zap.Logger, routerURL string, config Mess }, nil } -func (asc AzureStorageConnection) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubscription, error) { +func (asc AzureStorageConnection) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Subscription, error) { asc.logger.Info("subscribing to Azure storage queue", zap.String("queue", trigger.Spec.Topic)) if trigger.Spec.FunctionReference.Type != fv1.FunctionReferenceTypeFunctionName { @@ -224,7 +230,7 @@ func (asc AzureStorageConnection) subscribe(trigger *fv1.MessageQueueTrigger) (m return subscription, nil } -func (asc AzureStorageConnection) unsubscribe(subscription messageQueueSubscription) error { +func (asc AzureStorageConnection) Unsubscribe(subscription messageQueue.Subscription) error { sub := subscription.(*AzureQueueSubscription) asc.logger.Info("unsubscribing from Azure storage queue", zap.String("queue", sub.queueName)) @@ -394,3 +400,7 @@ func invokeTriggeredFunction(conn AzureStorageConnection, sub *AzureQueueSubscri return } } + +func IsTopicValid(topic string) bool { + return len(topic) >= 3 && len(topic) <= 63 && validAzureQueueName.MatchString(topic) +} diff --git a/pkg/mqtrigger/messageQueue/asq_test.go b/pkg/mqtrigger/messageQueue/azurequeuestorage/asq_test.go similarity index 96% rename from pkg/mqtrigger/messageQueue/asq_test.go rename to pkg/mqtrigger/messageQueue/azurequeuestorage/asq_test.go index 8018181f..fa4af81c 100644 --- a/pkg/mqtrigger/messageQueue/asq_test.go +++ b/pkg/mqtrigger/messageQueue/azurequeuestorage/asq_test.go @@ -13,7 +13,7 @@ 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 messageQueue +package azurequeuestorage import ( "fmt" @@ -32,6 +32,7 @@ import ( metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/mqtrigger/messageQueue" ) const ( @@ -112,7 +113,7 @@ func TestNewStorageConnectionMissingAccountName(t *testing.T) { logger, err := zap.NewDevelopment() panicIf(err) - connection, err := newAzureStorageConnection(logger, DummyRouterURL, MessageQueueConfig{ + connection, err := New(logger, DummyRouterURL, messageQueue.Config{ MQType: fv1.MessageQueueTypeASQ, Url: "", }) @@ -125,7 +126,7 @@ func TestNewStorageConnectionMissingAccessKey(t *testing.T) { panicIf(err) _ = os.Setenv("AZURE_STORAGE_ACCOUNT_NAME", "accountname") - connection, err := newAzureStorageConnection(logger, DummyRouterURL, MessageQueueConfig{ + connection, err := New(logger, DummyRouterURL, messageQueue.Config{ MQType: fv1.MessageQueueTypeASQ, Url: "", }) @@ -140,7 +141,7 @@ func TestNewStorageConnection(t *testing.T) { _ = os.Setenv("AZURE_STORAGE_ACCOUNT_NAME", "accountname") _ = os.Setenv("AZURE_STORAGE_ACCOUNT_KEY", "bm90IGEga2V5") - connection, err := newAzureStorageConnection(logger, DummyRouterURL, MessageQueueConfig{ + connection, err := New(logger, DummyRouterURL, messageQueue.Config{ MQType: "azure-storage-queue", Url: "", }) @@ -301,7 +302,7 @@ func TestAzureStorageQueuePoisonMessage(t *testing.T) { service: service, httpClient: httpClient, } - subscription, err := connection.subscribe(&fv1.MessageQueueTrigger{ + subscription, err := connection.Subscribe(&fv1.MessageQueueTrigger{ ObjectMeta: metav1.ObjectMeta{ Name: TriggerName, Namespace: metav1.NamespaceDefault, @@ -319,7 +320,7 @@ func TestAzureStorageQueuePoisonMessage(t *testing.T) { require.NoError(t, err) require.NotNil(t, subscription) - connection.unsubscribe(subscription) + connection.Unsubscribe(subscription) mock.AssertExpectationsForObjects(t, httpClient, message, poisonMessage, queue, poisonQueue, service) } @@ -449,7 +450,7 @@ func runAzureStorageQueueTest(t *testing.T, count int, output bool) { service: service, httpClient: httpClient, } - subscription, err := connection.subscribe(&fv1.MessageQueueTrigger{ + subscription, err := connection.Subscribe(&fv1.MessageQueueTrigger{ ObjectMeta: metav1.ObjectMeta{ Name: TriggerName, Namespace: metav1.NamespaceDefault, @@ -468,7 +469,7 @@ func runAzureStorageQueueTest(t *testing.T, count int, output bool) { require.NoError(t, err) require.NotNil(t, subscription) - connection.unsubscribe(subscription) + connection.Unsubscribe(subscription) mock.AssertExpectationsForObjects(t, httpClient, message, outputMessage, queue, outputQueue, service) } diff --git a/pkg/mqtrigger/messageQueue/kafka.go b/pkg/mqtrigger/messageQueue/kafka/kafka.go similarity index 91% rename from pkg/mqtrigger/messageQueue/kafka.go rename to pkg/mqtrigger/messageQueue/kafka/kafka.go index b0d3e458..d54d5ba7 100644 --- a/pkg/mqtrigger/messageQueue/kafka.go +++ b/pkg/mqtrigger/messageQueue/kafka/kafka.go @@ -14,7 +14,7 @@ See the License for the specific language governing permissions and limitations under the License. */ -package messageQueue +package kafka import ( "crypto/tls" @@ -23,6 +23,7 @@ import ( "io/ioutil" "net/http" "os" + "regexp" "strconv" "strings" @@ -32,9 +33,15 @@ import ( "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/mqtrigger/messageQueue" "github.com/fission/fission/pkg/utils" ) +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]$`) +) + type ( Kafka struct { logger *zap.Logger @@ -46,7 +53,7 @@ type ( } ) -func makeKafkaMessageQueue(logger *zap.Logger, routerUrl string, mqCfg MessageQueueConfig) (MessageQueue, error) { +func New(logger *zap.Logger, routerUrl string, mqCfg messageQueue.Config) (messageQueue.MessageQueue, error) { if len(routerUrl) == 0 || len(mqCfg.Url) == 0 { return nil, errors.New("the router URL or MQ URL is empty") } @@ -88,11 +95,7 @@ func makeKafkaMessageQueue(logger *zap.Logger, routerUrl string, mqCfg MessageQu return kafka, nil } -func isTopicValidForKafka(topic string) bool { - return true -} - -func (kafka Kafka) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubscription, error) { +func (kafka Kafka) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Subscription, error) { kafka.logger.Info("inside kakfa subscribe", zap.Any("trigger", trigger)) kafka.logger.Info("brokers set", zap.Strings("brokers", kafka.brokers)) @@ -194,7 +197,7 @@ func (kafka Kafka) getTLSConfig() (*tls.Config, error) { return &tlsConfig, nil } -func (kafka Kafka) unsubscribe(subscription messageQueueSubscription) error { +func (kafka Kafka) Unsubscribe(subscription messageQueue.Subscription) error { return subscription.(*cluster.Consumer).Close() } @@ -336,3 +339,20 @@ func errorHandler(logger *zap.Logger, trigger *fv1.MessageQueueTrigger, producer zap.String("message", err.Error()), zap.String("trigger", trigger.ObjectMeta.Name), zap.String("function_url", funcUrl)) } } + +// 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 + } + if topic == "." || topic == ".." { + return false + } + if len(topic) > 249 { + return false + } + if !validKafkaTopicName.MatchString(topic) { + return false + } + return true +} diff --git a/pkg/mqtrigger/messageQueue/messagequeue.go b/pkg/mqtrigger/messageQueue/messagequeue.go new file mode 100644 index 00000000..ce9241ef --- /dev/null +++ b/pkg/mqtrigger/messageQueue/messagequeue.go @@ -0,0 +1,36 @@ +/* +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 messageQueue + +import ( + fv1 "github.com/fission/fission/pkg/apis/core/v1" +) + +type ( + Subscription interface{} + + Config struct { + MQType string + Url string + Secrets map[string][]byte + } + + MessageQueue interface { + Subscribe(trigger *fv1.MessageQueueTrigger) (Subscription, error) + Unsubscribe(triggerSub Subscription) error + } +) diff --git a/pkg/mqtrigger/messageQueue/nats.go b/pkg/mqtrigger/messageQueue/nats/nats.go similarity index 94% rename from pkg/mqtrigger/messageQueue/nats.go rename to pkg/mqtrigger/messageQueue/nats/nats.go index d02564e5..52aadd4d 100644 --- a/pkg/mqtrigger/messageQueue/nats.go +++ b/pkg/mqtrigger/messageQueue/nats/nats.go @@ -14,7 +14,7 @@ See the License for the specific language governing permissions and limitations under the License. */ -package messageQueue +package nats import ( "bytes" @@ -28,6 +28,7 @@ import ( "go.uber.org/zap" fv1 "github.com/fission/fission/pkg/apis/core/v1" + "github.com/fission/fission/pkg/mqtrigger/messageQueue" "github.com/fission/fission/pkg/utils" ) @@ -46,7 +47,7 @@ type ( } ) -func makeNatsMessageQueue(logger *zap.Logger, routerUrl string, mqCfg MessageQueueConfig) (MessageQueue, error) { +func New(logger *zap.Logger, routerUrl string, mqCfg messageQueue.Config) (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. @@ -68,10 +69,10 @@ func makeNatsMessageQueue(logger *zap.Logger, routerUrl string, mqCfg MessageQue return nats, nil } -func (nats Nats) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubscription, error) { +func (nats Nats) Subscribe(trigger *fv1.MessageQueueTrigger) (messageQueue.Subscription, error) { subj := trigger.Spec.Topic - if !isTopicValidForNats(subj) { + if !IsTopicValid(subj) { return nil, fmt.Errorf("not a valid topic: %q", trigger.Spec.Topic) } @@ -92,15 +93,10 @@ func (nats Nats) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubscr return sub, nil } -func (nats Nats) unsubscribe(subscription messageQueueSubscription) error { +func (nats Nats) Unsubscribe(subscription messageQueue.Subscription) error { return subscription.(ns.Subscription).Close() } -func isTopicValidForNats(topic string) bool { - // nats-streaming does not support wildcard channel. - return nsUtil.IsChannelNameValid(topic, false) -} - func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) { return func(msg *ns.Msg) { @@ -212,5 +208,9 @@ func msgHandler(nats *Nats, trigger *fv1.MessageQueueTrigger) func(*ns.Msg) { } } } - +} + +func IsTopicValid(topic string) bool { + // nats-streaming does not support wildcard channel. + return nsUtil.IsChannelNameValid(topic, false) } diff --git a/pkg/mqtrigger/messageQueue/messageQueue.go b/pkg/mqtrigger/mqtmanager.go similarity index 78% rename from pkg/mqtrigger/messageQueue/messageQueue.go rename to pkg/mqtrigger/mqtmanager.go index e56968c4..0ec670f3 100644 --- a/pkg/mqtrigger/messageQueue/messageQueue.go +++ b/pkg/mqtrigger/mqtmanager.go @@ -14,19 +14,22 @@ See the License for the specific language governing permissions and limitations under the License. */ -package messageQueue +package mqtrigger import ( "errors" - "fmt" "time" - "github.com/fission/fission/pkg/utils" "go.uber.org/zap" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" 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" ) const ( @@ -36,32 +39,19 @@ const ( ) type ( - messageQueueSubscription interface{} - requestType int - MessageQueueConfig struct { - MQType string - Url string - Secrets map[string][]byte - } - - MessageQueue interface { - subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubscription, error) - unsubscribe(triggerSub messageQueueSubscription) error - } - MessageQueueTriggerManager struct { logger *zap.Logger reqChan chan request triggers map[string]*triggerSubscription fissionClient *crd.FissionClient - messageQueue MessageQueue + messageQueue messageQueue.MessageQueue } triggerSubscription struct { trigger fv1.MessageQueueTrigger - subscription messageQueueSubscription + subscription messageQueue.Subscription } request struct { @@ -75,35 +65,23 @@ type ( } ) -func MakeMessageQueueTriggerManager(logger *zap.Logger, fissionClient *crd.FissionClient, routerUrl string, mqConfig MessageQueueConfig) *MessageQueueTriggerManager { - var messageQueue MessageQueue - var err error - +func MakeMessageQueueTriggerManager(logger *zap.Logger, + fissionClient *crd.FissionClient, 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, } - switch mqConfig.MQType { - case fv1.MessageQueueTypeNats: - messageQueue, err = makeNatsMessageQueue(logger, routerUrl, mqConfig) - case fv1.MessageQueueTypeASQ: - messageQueue, err = newAzureStorageConnection(logger, routerUrl, mqConfig) - case fv1.MessageQueueTypeKafka: - messageQueue, err = makeKafkaMessageQueue(logger, routerUrl, mqConfig) - default: - err = fmt.Errorf("no supported message queue type found for %q", mqConfig.MQType) - } - if err != nil { - logger.Fatal("failed to connect to remote message queue server", zap.Error(err)) - } - mqTriggerMgr.messageQueue = messageQueue - go mqTriggerMgr.service() - go mqTriggerMgr.syncTriggers() return &mqTriggerMgr } +func (mqt *MessageQueueTriggerManager) Run() { + go mqt.service() + go mqt.syncTriggers() +} + func (mqt *MessageQueueTriggerManager) service() { for { req := <-mqt.reqChan @@ -189,7 +167,7 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() { } // actually subscribe using the message queue client impl - sub, err := mqt.messageQueue.subscribe(trigger) + sub, err := mqt.messageQueue.Subscribe(trigger) if err != nil { mqt.logger.Warn("failed to subscribe to message queue trigger", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name)) continue @@ -214,7 +192,7 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() { if _, ok := newTriggerMap[key]; ok { continue } - err := mqt.messageQueue.unsubscribe(triggerSub.subscription) + err := mqt.messageQueue.Unsubscribe(triggerSub.subscription) if err != nil { mqt.logger.Warn("failed to unsubscribe from message queue trigger", zap.Error(err), zap.String("trigger_name", triggerSub.trigger.ObjectMeta.Name)) continue @@ -228,12 +206,14 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() { } } -func IsTopicValid(mqType string, topic string) bool { +func IsTopicValid(mqType fv1.MessageQueueType, topic string) bool { switch mqType { case fv1.MessageQueueTypeNats: - return isTopicValidForNats(topic) + return nats.IsTopicValid(topic) + case fv1.MessageQueueTypeASQ: + return azurequeuestorage.IsTopicValid(topic) case fv1.MessageQueueTypeKafka: - return isTopicValidForKafka(topic) + return kafka.IsTopicValid(topic) } return false }