Reorganize message queue trigger directory structure (#1531)

This commit is contained in:
Ta-Ching Chen
2020-02-12 15:40:00 +08:00
committed by GitHub
parent d2e9364e5e
commit 75321d306a
8 changed files with 156 additions and 82 deletions
@@ -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)
}
@@ -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)
}
@@ -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
}
-239
View File
@@ -1,239 +0,0 @@
/*
Copyright 2016 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 (
"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"
)
const (
ADD_TRIGGER requestType = iota
DELETE_TRIGGER
GET_ALL_TRIGGERS
)
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
}
triggerSubscription struct {
trigger fv1.MessageQueueTrigger
subscription messageQueueSubscription
}
request struct {
requestType
triggerSub *triggerSubscription
respChan chan response
}
response struct {
err error
triggers *map[string]*triggerSubscription
}
)
func MakeMessageQueueTriggerManager(logger *zap.Logger, fissionClient *crd.FissionClient, routerUrl string, mqConfig MessageQueueConfig) *MessageQueueTriggerManager {
var messageQueue MessageQueue
var err error
mqTriggerMgr := MessageQueueTriggerManager{
logger: logger.Named("message_queue_trigger_manager"),
reqChan: make(chan request),
triggers: make(map[string]*triggerSubscription),
fissionClient: fissionClient,
}
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) service() {
for {
req := <-mqt.reqChan
switch req.requestType {
case ADD_TRIGGER:
var err error
k := crd.CacheKey(&req.triggerSub.trigger.ObjectMeta)
if _, ok := mqt.triggers[k]; ok {
err = errors.New("trigger already exists")
} else {
mqt.triggers[k] = req.triggerSub
}
req.respChan <- response{err: err}
case GET_ALL_TRIGGERS:
copyTriggers := make(map[string]*triggerSubscription)
for key, val := range mqt.triggers {
copyTriggers[key] = val
}
req.respChan <- response{triggers: &copyTriggers}
case DELETE_TRIGGER:
delete(mqt.triggers, crd.CacheKey(&req.triggerSub.trigger.ObjectMeta))
}
}
}
func (mqt *MessageQueueTriggerManager) addTrigger(triggerSub *triggerSubscription) error {
respChan := make(chan response)
mqt.reqChan <- request{
requestType: ADD_TRIGGER,
triggerSub: triggerSub,
respChan: respChan,
}
r := <-respChan
return r.err
}
func (mqt *MessageQueueTriggerManager) getAllTriggers() *map[string]*triggerSubscription {
respChan := make(chan response)
mqt.reqChan <- request{
requestType: GET_ALL_TRIGGERS,
respChan: respChan,
}
r := <-respChan
return r.triggers
}
func (mqt *MessageQueueTriggerManager) delTrigger(m *metav1.ObjectMeta) {
mqt.reqChan <- request{
requestType: DELETE_TRIGGER,
triggerSub: &triggerSubscription{
trigger: fv1.MessageQueueTrigger{
ObjectMeta: *m,
},
},
}
}
func (mqt *MessageQueueTriggerManager) syncTriggers() {
for {
// get new set of triggers
newTriggers, err := mqt.fissionClient.CoreV1().MessageQueueTriggers(metav1.NamespaceAll).List(metav1.ListOptions{})
if err != nil {
if utils.IsNetworkError(err) {
mqt.logger.Error("encountered network error, will retry", zap.Error(err))
time.Sleep(5 * time.Second)
continue
}
mqt.logger.Fatal("failed to read message queue trigger list", zap.Error(err))
}
newTriggerMap := make(map[string]*fv1.MessageQueueTrigger)
for index := range newTriggers.Items {
newTrigger := &newTriggers.Items[index]
newTriggerMap[crd.CacheKey(&newTrigger.ObjectMeta)] = newTrigger
}
// get current set of triggers
currentTriggers := mqt.getAllTriggers()
// register new triggers
for key, trigger := range newTriggerMap {
if _, ok := (*currentTriggers)[key]; ok {
continue
}
// actually subscribe using the message queue client impl
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
}
triggerSub := triggerSubscription{
trigger: *trigger,
subscription: sub,
}
// add to our list
err = mqt.addTrigger(&triggerSub)
if err != nil {
mqt.logger.Fatal("adding message queue trigger failed", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name))
}
mqt.logger.Info("message queue trigger created", zap.String("trigger_name", trigger.ObjectMeta.Name))
}
// remove old triggers
for key, triggerSub := range *currentTriggers {
if _, ok := newTriggerMap[key]; ok {
continue
}
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
}
mqt.delTrigger(&triggerSub.trigger.ObjectMeta)
mqt.logger.Info("message queue trigger deleted", zap.String("trigger_name", triggerSub.trigger.ObjectMeta.Name))
}
// TODO replace with a watch
time.Sleep(3 * time.Second)
}
}
func IsTopicValid(mqType string, topic string) bool {
switch mqType {
case fv1.MessageQueueTypeNats:
return isTopicValidForNats(topic)
case fv1.MessageQueueTypeKafka:
return isTopicValidForKafka(topic)
}
return false
}
@@ -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
}
)
@@ -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)
}