Add message queue service factory (#1537)
This commit is contained in:
@@ -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) {
|
||||
|
||||
@@ -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 (
|
||||
|
||||
@@ -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,19 +494,17 @@ 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"))
|
||||
}
|
||||
|
||||
if !IsTopicValid(spec.MessageQueueType, spec.Topic) {
|
||||
} else {
|
||||
if !validator.IsValidTopic((string)(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) {
|
||||
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()
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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"}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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"
|
||||
)
|
||||
|
||||
@@ -46,6 +43,7 @@ type (
|
||||
reqChan chan request
|
||||
triggers map[string]*triggerSubscription
|
||||
fissionClient *crd.FissionClient
|
||||
messageQueueType fv1.MessageQueueType
|
||||
messageQueue messageQueue.MessageQueue
|
||||
}
|
||||
|
||||
@@ -66,12 +64,13 @@ 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,
|
||||
messageQueueType: mqType,
|
||||
messageQueue: messageQueue,
|
||||
}
|
||||
return &mqTriggerMgr
|
||||
@@ -154,8 +153,10 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() {
|
||||
newTriggerMap := make(map[string]*fv1.MessageQueueTrigger)
|
||||
for index := range newTriggers.Items {
|
||||
newTrigger := &newTriggers.Items[index]
|
||||
if newTrigger.Spec.MessageQueueType == mqt.messageQueueType {
|
||||
newTriggerMap[crd.CacheKey(&newTrigger.ObjectMeta)] = newTrigger
|
||||
}
|
||||
}
|
||||
|
||||
// get current set of triggers
|
||||
currentTriggers := mqt.getAllTriggers()
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
}
|
||||
Reference in New Issue
Block a user