From 96cbee185da65e81c440b93a97df203d5cd695d4 Mon Sep 17 00:00:00 2001 From: Vishal Date: Mon, 1 Oct 2018 18:03:31 +0530 Subject: [PATCH] Kafka integration (#831) Kafka integration with Fission enables invoking a function when a message arrives in a Kafka topic --- charts/README.md | 15 +- charts/fission-all/templates/deployment.yaml | 73 ++++++- charts/fission-all/templates/secrets.yaml | 2 +- charts/fission-all/templates/svc.yaml | 5 +- charts/fission-all/values.yaml | 16 +- fission/mqtrigger.go | 8 +- glide.lock | 24 ++- glide.yaml | 4 + mqtrigger/messageQueue/kafka.go | 197 +++++++++++++++++++ mqtrigger/messageQueue/kafka_test.go | 119 +++++++++++ mqtrigger/messageQueue/messageQueue.go | 13 ++ pkg/apis/fission.io/v1/const.go | 5 +- pkg/apis/fission.io/v1/validation.go | 22 ++- test/test_utils.sh | 2 +- test/tests/test_package_command.sh | 2 + types.go | 5 +- 16 files changed, 473 insertions(+), 39 deletions(-) create mode 100644 mqtrigger/messageQueue/kafka.go create mode 100644 mqtrigger/messageQueue/kafka_test.go diff --git a/charts/README.md b/charts/README.md index d6cde451..27726aaf 100644 --- a/charts/README.md +++ b/charts/README.md @@ -56,11 +56,16 @@ The following table lists the configurable parameters of the Fission chart and t | `logger.influxdbAdmin` | Log database admin username | `admin` | | `logger.fluentdImage` | Logger fluentd image | `fission/fluentd` | | `fissionUiImage` | Fission ui image | `fission/fission-ui:0.1.0` | -| `messageQueue` | Message queue type | `nats-streaming` | -| `nats.authToken` | Nats streaming auth token | `defaultFissionAuthToken` | -| `nats.clusterID` | Nats streaming clusterID | `fissionMQTrigger` | -| `azureStorageQueue.accountName` | Azure storage account name | None (required if `messageQueue` is `azure-storage-queue`) | -| `azureStorageQueue.key` | Azure storage access key | None (required if `messageQueue` is `azure-storage-queue`) | +| `nats.enabled` | Nats streaming enabled | `true` | +| `nats.authToken` | Nats streaming auth token | `defaultFissionAuthToken`(required if `nats.enabled` is `true`) | +| `nats.clusterID` | Nats streaming clusterID | `fissionMQTrigger`(required if `nats.enabled` is `true`) | +| `azureStorageQueue.enabled` * | Azure storage account name | false | +| `azureStorageQueue.accountName` | Azure storage account name | None (required if `azureStorageQueue.enabled` is `true`) | +| `azureStorageQueue.key` | Azure storage access key | None (required if `azureStorageQueue.enabled` is `true`) | +| `kafka.enabled` * | Kafka trigger enabled | `false` | +| `kafka.brokers` | Kafka brokers uri | `broker.kafka:9092` (required if `kafka.enabled` is `true`) | + +* - Please note that deploying of Azure Storage Queue or Kafka is not done by Fission chart and you will have to explicitly deploy them. Specify each parameter using the `--set key=value[,key=value]` argument to `helm install`. For example, diff --git a/charts/fission-all/templates/deployment.yaml b/charts/fission-all/templates/deployment.yaml index dd7105d1..d14f69d5 100644 --- a/charts/fission-all/templates/deployment.yaml +++ b/charts/fission-all/templates/deployment.yaml @@ -342,7 +342,7 @@ metadata: svc: influxdb chart: "{{ .Chart.Name }}-{{ .Chart.Version }}" spec: - type: ClusterIP + type: ClusterIP ports: - port: 8086 targetPort: 8086 @@ -531,7 +531,7 @@ spec: # serviceAccount: fission-svc --- -{{- if eq .Values.messageQueue.type "nats-streaming" }} +{{- if .Values.nats.enabled }} apiVersion: extensions/v1beta1 kind: Deployment metadata: @@ -553,17 +553,16 @@ spec: "--auth", "{{ .Values.nats.authToken }}", "--max_channels", "0" ] + ports: - containerPort: 4222 hostPort: 4222 protocol: TCP -{{- end }} - --- apiVersion: extensions/v1beta1 kind: Deployment metadata: - name: mqtrigger + name: mqtrigger-nats-streaming labels: chart: "{{ .Chart.Name }}-{{ .Chart.Version }}" spec: @@ -572,6 +571,7 @@ spec: metadata: labels: svc: mqtrigger + messagequeue: nats-streaming spec: containers: - name: mqtrigger @@ -581,11 +581,65 @@ spec: args: ["--mqt", "--routerUrl", "http://router.{{ .Release.Namespace }}"] env: - name: MESSAGE_QUEUE_TYPE - value: {{ .Values.messageQueue.type }} - {{- if eq .Values.messageQueue.type "nats-streaming" }} + value: nats-streaming - name: MESSAGE_QUEUE_URL value: nats://{{ .Values.nats.authToken }}@nats-streaming:4222 - {{- else if eq .Values.messageQueue.type "azure-storage-queue" }} + serviceAccount: fission-svc +{{- end }} +--- +{{- if .Values.kafka.enabled }} +apiVersion: extensions/v1beta1 +kind: Deployment +metadata: + name: mqtrigger-kafka + labels: + chart: "{{ .Chart.Name }}-{{ .Chart.Version }}" +spec: + replicas: 1 + template: + metadata: + labels: + svc: mqtrigger + messagequeue: kafka + spec: + containers: + - name: mqtrigger + image: "{{ .Values.image }}:{{ .Values.imageTag }}" + imagePullPolicy: {{ .Values.pullPolicy }} + command: ["/fission-bundle"] + args: ["--mqt"] + env: + - name: MESSAGE_QUEUE_TYPE + value: kafka + - name: MESSAGE_QUEUE_URL + value: "{{.Values.kafka.brokers}}" + serviceAccount: fission-svc +{{- end }} +--- +{{- if .Values.azureStorageQueue.enabled }} +apiVersion: extensions/v1beta1 +kind: Deployment +metadata: + name: mqtrigger-azure-storage-queue + labels: + chart: "{{ .Chart.Name }}-{{ .Chart.Version }}" +spec: + replicas: 1 + template: + metadata: + labels: + svc: mqtrigger + messagequeue: azure-storage-queue + spec: + containers: + - name: mqtrigger + image: "{{ .Values.image }}:{{ .Values.imageTag }}" + imagePullPolicy: {{ .Values.pullPolicy }} + command: ["/fission-bundle"] + args: ["--mqt", "--routerUrl", "http://router.{{ .Release.Namespace }}"] + env: + - name: MESSAGE_QUEUE_TYPE + value: azure-storage-queue - name: AZURE_STORAGE_ACCOUNT_NAME value: {{ required "An Azure storage account name is required." .Values.azureStorageQueue.accountName }} - name: AZURE_STORAGE_ACCOUNT_KEY @@ -593,9 +647,8 @@ spec: secretKeyRef: name: azure-storage-account-key key: key - {{- end }} serviceAccount: fission-svc - +{{- end }} --- apiVersion: extensions/v1beta1 kind: Deployment diff --git a/charts/fission-all/templates/secrets.yaml b/charts/fission-all/templates/secrets.yaml index 00972cba..c4abc8e3 100644 --- a/charts/fission-all/templates/secrets.yaml +++ b/charts/fission-all/templates/secrets.yaml @@ -10,7 +10,7 @@ data: password: {{ randAlphaNum 20 | b64enc | quote }} --- -{{- if eq .Values.messageQueue.type "azure-storage-queue" }} +{{- if .Values.azureStorageQueue.enabled }} apiVersion: v1 kind: Secret metadata: diff --git a/charts/fission-all/templates/svc.yaml b/charts/fission-all/templates/svc.yaml index 0dbfbbce..a9b06dd2 100644 --- a/charts/fission-all/templates/svc.yaml +++ b/charts/fission-all/templates/svc.yaml @@ -38,7 +38,7 @@ spec: svc: controller --- -{{- if eq .Values.messageQueue.type "nats-streaming" }} +{{- if .Values.nats.enabled }} apiVersion: v1 kind: Service metadata: @@ -57,7 +57,6 @@ spec: selector: svc: nats-streaming {{- end }} - --- apiVersion: v1 kind: Service @@ -73,4 +72,4 @@ spec: - port: 80 targetPort: 8000 selector: - svc: storagesvc + svc: storagesvc \ No newline at end of file diff --git a/charts/fission-all/values.yaml b/charts/fission-all/values.yaml index df896aba..bea6cdd6 100644 --- a/charts/fission-all/values.yaml +++ b/charts/fission-all/values.yaml @@ -53,21 +53,23 @@ logger: fluentdImage: fission/fluentd fluentdImageTag: 0.10.0 -## Type of Queue you would like to use -## currently supports nats-streaming, azure-storage-queue -messageQueue: - type: nats-streaming - ## Message queue trigger config -### NATS Streaming +### NATS Streaming, enabled by default nats: + enabled: true authToken: "defaultFissionAuthToken" clusterID: "fissionMQTrigger" -## Required if messageQueue type is azure-storage-queue +## Azure-storage-queue: enable and configure the details azureStorageQueue: + enabled: false key: "" accountName: "" + +## Kafka: enable and configure the details +kafka: + enabled: false + brokers: 'broker.kafka:9092' ## Persist data to a persistent volume. persistence: diff --git a/fission/mqtrigger.go b/fission/mqtrigger.go index a3192885..ea47ce63 100644 --- a/fission/mqtrigger.go +++ b/fission/mqtrigger.go @@ -53,14 +53,18 @@ func mqtCreate(c *cli.Context) error { mqType = fission.MessageQueueTypeNats case fission.MessageQueueTypeASQ: mqType = fission.MessageQueueTypeASQ + case fission.MessageQueueTypeKafka: + mqType = fission.MessageQueueTypeKafka + default: - log.Fatal("Unknown message queue type, currently only \"nats-streaming, azure-storage-queue \" is supported") + log.Fatal("Unknown message queue type, currently only \"nats-streaming, azure-storage-queue, kafka \" is supported") + } // TODO: check topic availability topic := c.String("topic") if len(topic) == 0 { - log.Fatal("Listen topic cannot be empty") + log.Fatal("Topic cannot be empty") } respTopic := c.String("resptopic") diff --git a/glide.lock b/glide.lock index d33ddf37..1aa3b3a0 100644 --- a/glide.lock +++ b/glide.lock @@ -1,5 +1,5 @@ -hash: b26f975bf145378e939db833fa4b6a142fa04b884463d8173622b30f95191d42 -updated: 2018-08-13T13:12:57.58715-07:00 +hash: 07bc3f35b3b63ff6f15edf3b8724eb61bd6427572322fc171efbe18cec91cd9e +updated: 2018-09-17T14:46:13.995469387+05:30 imports: - name: cloud.google.com/go version: 3b1ae45394a234c385be014e9a488f2bb6eef821 @@ -21,12 +21,14 @@ imports: version: 3ac7bf7a47d159a033b107610db8a1b6575507a4 subpackages: - quantile +- name: github.com/bsm/sarama-cluster + version: c618e605e15c0d7535f6c96ff8efbb0dba4fd66c - name: github.com/coreos/etcd version: f87b566248bb0713a56dc55bc545aa5aad17ace0 subpackages: - client - name: github.com/davecgh/go-spew - version: 346938d642f2ec3594ed81d874461961cd0faa76 + version: 8991bc29aa16c548c550c7ff78260e27b9ab7c73 subpackages: - spew - name: github.com/dchest/uniuri @@ -49,6 +51,14 @@ imports: - internal/prefix - name: github.com/dustin/go-humanize version: 9f541cc9db5d55bce703bd99987c9d5cb8eea45e +- name: github.com/eapache/go-resiliency + version: ea41b0fad31007accc7f806884dcdf3da98b79ce + subpackages: + - breaker +- name: github.com/eapache/go-xerial-snappy + version: 040cc1a32f578808623071247fdbd5cc43f37f5f +- name: github.com/eapache/queue + version: 093482f3f8ce946c05bcba64badd2c82369e084d - name: github.com/fsnotify/fsnotify version: c2828203cd70a50dcccfb2761f8b1f8ceef9a8e9 - name: github.com/ghodss/yaml @@ -137,7 +147,7 @@ imports: - name: github.com/modern-go/reflect2 version: 05fbef0ca5da472bbf96c9322b84a53edc03c9fd - name: github.com/nats-io/go-nats - version: 2485387d6ede89c1c8c3445cd1935fd41a2e9ee9 + version: fb0396ee0bdb8018b0fef30d6d1de798ce99cd05 subpackages: - encoders/builtin - util @@ -146,7 +156,7 @@ imports: subpackages: - pb - name: github.com/nats-io/nats-streaming-server - version: 63e2c334b66dba3edade0047625a25c0d8f18f80 + version: 8910c0c347bc51cc87227aeb27bef19d409bf5c2 subpackages: - spb - util @@ -179,10 +189,14 @@ imports: version: 65c1f6f8f0fc1e2185eb9863a3bc751496404259 subpackages: - xfs +- name: github.com/rcrowley/go-metrics + version: e2704e165165ec55d062f5919b4b29494e9fa790 - name: github.com/robfig/cron version: b41be1df696709bb6395fe435af20370037c0b4c - name: github.com/satori/go.uuid version: f58768cc1a7a7e77a3bd49e98cdd21419399b6a3 +- name: github.com/Shopify/sarama + version: a6144ae922fd99dd0ea5046c8137acfb7fab0914 - name: github.com/sirupsen/logrus version: 68cec9f21fbf3ea8d8f98c044bc6ce05f17b267a - name: github.com/spf13/pflag diff --git a/glide.yaml b/glide.yaml index 045f572c..2aeee8f2 100644 --- a/glide.yaml +++ b/glide.yaml @@ -70,3 +70,7 @@ import: - package: github.com/dustin/go-humanize - package: github.com/golang/protobuf/proto version: v1.1.0 +- package: github.com/bsm/sarama-cluster + version: ^2.1.11 +- package: github.com/Shopify/sarama + version: ^1.15.0 diff --git a/mqtrigger/messageQueue/kafka.go b/mqtrigger/messageQueue/kafka.go new file mode 100644 index 00000000..3fd0ab92 --- /dev/null +++ b/mqtrigger/messageQueue/kafka.go @@ -0,0 +1,197 @@ +/* +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" + "io/ioutil" + "net/http" + "strings" + + sarama "github.com/Shopify/sarama" + cluster "github.com/bsm/sarama-cluster" + log "github.com/sirupsen/logrus" + + "github.com/fission/fission" + "github.com/fission/fission/crd" +) + +type ( + Kafka struct { + routerUrl string + brokers []string + } +) + +func makeKafkaMessageQueue(routerUrl string, mqCfg MessageQueueConfig) (MessageQueue, error) { + if len(routerUrl) == 0 || len(mqCfg.Url) == 0 { + return nil, errors.New("The router URL or MQ URL is empty") + } + kafka := Kafka{ + routerUrl: routerUrl, + brokers: strings.Split(mqCfg.Url, ","), + } + log.Infof("Created Queue ", kafka) + return kafka, nil +} + +func isTopicValidForKafka(topic string) bool { + return true +} + +func (kafka Kafka) subscribe(trigger *crd.MessageQueueTrigger) (messageQueueSubscription, error) { + log.Infof("Inside kakfa subscribe", trigger) + log.Infof("brokers set to ", kafka.brokers) + + // Create new consumer + consumerConfig := cluster.NewConfig() + consumerConfig.Consumer.Return.Errors = true + consumerConfig.Group.Return.Notifications = true + consumer, err := cluster.NewConsumer(kafka.brokers, string(trigger.Metadata.UID), []string{trigger.Spec.Topic}, consumerConfig) + log.Infof("Created a new consumer ", consumer) + if err != nil { + panic(err) + } + + // Create new producer + producerConfig := sarama.NewConfig() + producerConfig.Producer.RequiredAcks = sarama.WaitForAll + producerConfig.Producer.Retry.Max = 10 + producerConfig.Producer.Return.Successes = true + producer, err := sarama.NewSyncProducer(kafka.brokers, producerConfig) + log.Infof("Created a new producer ", producer) + if err != nil { + panic(err) + } + + // consume errors + go func() { + for err := range consumer.Errors() { + log.Printf("Error: %s\n", err.Error()) + } + }() + + // consume notifications + go func() { + for ntf := range consumer.Notifications() { + log.Printf("Rebalanced: %+v\n", ntf) + } + }() + + // consume messages + go func() { + for msg := range consumer.Messages() { + log.Infof("Calling message handler with value " + string(msg.Value[:])) + if kafkaMsgHandler(&kafka, producer, trigger, string(msg.Value[:])) { + consumer.MarkOffset(msg, "") // mark message as processed + } + } + }() + + return consumer, nil +} + +func (kafka Kafka) unsubscribe(subscription messageQueueSubscription) error { + return subscription.(*cluster.Consumer).Close() +} + +func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *crd.MessageQueueTrigger, value string) bool { + // Support other function ref types + if trigger.Spec.FunctionReference.Type != fission.FunctionReferenceTypeFunctionName { + log.Fatalf("Unsupported function reference type (%v) for trigger %v", + trigger.Spec.FunctionReference.Type, trigger.Metadata.Name) + } + + url := kafka.routerUrl + "/" + strings.TrimPrefix(fission.UrlForFunction(trigger.Spec.FunctionReference.Name, trigger.Metadata.Namespace), "/") + log.Printf("Making HTTP request to %v", url) + headers := map[string]string{ + "X-Fission-MQTrigger-Topic": trigger.Spec.Topic, + "X-Fission-MQTrigger-RespTopic": trigger.Spec.ResponseTopic, + "X-Fission-MQTrigger-ErrorTopic": trigger.Spec.ErrorTopic, + "Content-Type": trigger.Spec.ContentType, + } + // Create request + req, err := http.NewRequest("POST", url, strings.NewReader(value)) + if err != nil { + log.Warningf("Request creation failed: %v", url) + return false + } + + for k, v := range headers { + req.Header.Add(k, v) + } + // Make the request + var resp *http.Response + for attempt := 0; attempt <= trigger.Spec.MaxRetries; attempt++ { + // Make the request + resp, err = http.DefaultClient.Do(req) + if err != nil { + log.Error("Error invoking function for trigger %v: %v", trigger.Metadata.Name, err) + continue + } + if resp == nil { + continue + } + if err == nil && resp.StatusCode == http.StatusOK { + // Success, quit retrying + break + } + } + + if resp == nil { + log.Warning("Every retry failed; final retry gave empty response.") + return false + } + defer resp.Body.Close() + body, err := ioutil.ReadAll(resp.Body) + log.Infof("Got response " + string(body)) + if err != nil { + errorHandler(trigger, producer, fmt.Sprintf("Request body error: %v", string(body))) + return false + } + if resp.StatusCode != 200 { + errorHandler(trigger, producer, fmt.Sprintf("Request returned failure: %v", resp.StatusCode)) + return false + } + if len(trigger.Spec.ResponseTopic) > 0 { + _, _, err := producer.SendMessage(&sarama.ProducerMessage{ + Topic: trigger.Spec.ResponseTopic, + Value: sarama.StringEncoder(body), + }) + if err != nil { + log.Warningf("Failed to publish message to topic %s: %v", trigger.Spec.ResponseTopic, err) + return false + } + } + return true +} + +func errorHandler(trigger *crd.MessageQueueTrigger, producer sarama.SyncProducer, body string) { + if len(trigger.Spec.ErrorTopic) > 0 { + _, _, err := producer.SendMessage(&sarama.ProducerMessage{ + Topic: trigger.Spec.ErrorTopic, + Value: sarama.StringEncoder(body), + }) + if err != nil { + log.Warningf("Failed to publish message to error topic %s: %v", trigger.Spec.ErrorTopic, err) + return + } + } else { + log.Printf(body) + } +} diff --git a/mqtrigger/messageQueue/kafka_test.go b/mqtrigger/messageQueue/kafka_test.go new file mode 100644 index 00000000..967d4856 --- /dev/null +++ b/mqtrigger/messageQueue/kafka_test.go @@ -0,0 +1,119 @@ +/* +Copyright 2017 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 ( + "testing" + + cluster "github.com/bsm/sarama-cluster" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + "github.com/fission/fission" + "github.com/fission/fission/crd" +) + +const ( + DummyKafkaBrokers = "http://kafka-0.kafka:9092,http://kafka-1.kafka:9092" +) + +type kafkaClusterMock struct { + mock.Mock +} + +type Consumer struct { + mock.Mock +} + +func (m *kafkaClusterMock) NewConsumer(addrs []string, groupID string, topics []string, config *cluster.Config) (*Consumer, error) { + args := m.Called(addrs, groupID, topics, config) + res := args.Get(0).(*Consumer) + //err := args.Error(1) + return res, nil +} + +func TestKafkaMQConfigValid(t *testing.T) { + kafkaConfig, err := makeKafkaMessageQueue(DummyRouterURL, MessageQueueConfig{ + MQType: fission.MessageQueueTypeKafka, + Url: DummyKafkaBrokers, + }) + require.NotNil(t, kafkaConfig) + require.Nil(t, err) +} + +func TestKafkaMQConfigMissingBroker(t *testing.T) { + kafkaConfig, err := makeKafkaMessageQueue(DummyRouterURL, MessageQueueConfig{ + MQType: fission.MessageQueueTypeKafka, + Url: "", + }) + require.Nil(t, kafkaConfig) + require.Error(t, err, "The router URL or MQ URL is empty") +} + +func TestKafkaMQConfigMissingRouter(t *testing.T) { + kafkaConfig, err := makeKafkaMessageQueue("", MessageQueueConfig{ + MQType: fission.MessageQueueTypeKafka, + Url: DummyKafkaBrokers, + }) + require.Nil(t, kafkaConfig) + require.Error(t, err, "The router URL or MQ URL is empty") +} + +func TestKafkaMq(t *testing.T) { + // This is a WIP test and does not yet work correctly, hence skipping for now + t.SkipNow() + + const ( + TriggerName = "queuetrigger" + QueueName = "inputqueue" + MessageBody = "input" + FunctionName = "testfunc" + ContentType = "text/plain" + ) + + kafkaConfig, err := makeKafkaMessageQueue(DummyRouterURL, MessageQueueConfig{ + MQType: fission.MessageQueueTypeKafka, + Url: DummyKafkaBrokers, + }) + + require.NoError(t, err) + + consumer := new(kafkaClusterMock) + consumer.On( + "NewConsumer", + mock.AnythingOfType("test"), + ).Return( + &Consumer{}, + nil, + ).Once() + + kafkaConfig.subscribe(&crd.MessageQueueTrigger{ + Metadata: metav1.ObjectMeta{ + Name: TriggerName, + Namespace: metav1.NamespaceDefault, + }, + Spec: fission.MessageQueueTriggerSpec{ + FunctionReference: fission.FunctionReference{ + Type: fission.FunctionReferenceTypeFunctionName, + Name: FunctionName, + }, + MessageQueueType: fission.MessageQueueTypeASQ, + Topic: QueueName, + ContentType: ContentType, + }, + }) +} diff --git a/mqtrigger/messageQueue/messageQueue.go b/mqtrigger/messageQueue/messageQueue.go index 6e13e053..0dd254a7 100644 --- a/mqtrigger/messageQueue/messageQueue.go +++ b/mqtrigger/messageQueue/messageQueue.go @@ -25,6 +25,7 @@ import ( "github.com/fission/fission" "github.com/fission/fission/crd" + fv1 "github.com/fission/fission/pkg/apis/fission.io/v1" ) const ( @@ -86,6 +87,8 @@ func MakeMessageQueueTriggerManager(fissionClient *crd.FissionClient, routerUrl messageQueue, err = makeNatsMessageQueue(routerUrl, mqConfig) case fission.MessageQueueTypeASQ: messageQueue, err = newAzureStorageConnection(routerUrl, mqConfig) + case fission.MessageQueueTypeKafka: + messageQueue, err = makeKafkaMessageQueue(routerUrl, mqConfig) default: err = errors.New("No matched message queue type found") } @@ -221,3 +224,13 @@ func (mqt *MessageQueueTriggerManager) syncTriggers() { 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 +} diff --git a/pkg/apis/fission.io/v1/const.go b/pkg/apis/fission.io/v1/const.go index d0d97521..80a159d3 100644 --- a/pkg/apis/fission.io/v1/const.go +++ b/pkg/apis/fission.io/v1/const.go @@ -64,8 +64,9 @@ const ( ) const ( - MessageQueueTypeNats = "nats-streaming" - MessageQueueTypeASQ = "azure-storage-queue" + MessageQueueTypeNats = "nats-streaming" + MessageQueueTypeASQ = "azure-storage-queue" + MessageQueueTypeKafka = "kafka" ) const ( diff --git a/pkg/apis/fission.io/v1/validation.go b/pkg/apis/fission.io/v1/validation.go index f019a6ed..ad9d2e77 100644 --- a/pkg/apis/fission.io/v1/validation.go +++ b/pkg/apis/fission.io/v1/validation.go @@ -36,6 +36,7 @@ const ( var ( validAzureQueueName = regexp.MustCompile("^[a-z0-9][a-z0-9\\-]*[a-z0-9]$") + validKafkaTopicName = regexp.MustCompile("^[a-z0-9][a-z0-9\\-._]*[a-z0-9]$") ) type ( @@ -164,10 +165,29 @@ func IsTopicValid(mqType MessageQueueType, topic string) bool { 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 @@ -433,7 +453,7 @@ func (spec MessageQueueTriggerSpec) Validate() error { result = multierror.Append(result, spec.FunctionReference.Validate()) switch spec.MessageQueueType { - case MessageQueueTypeNats, MessageQueueTypeASQ: // no op + case MessageQueueTypeNats, MessageQueueTypeASQ, MessageQueueTypeKafka: // no op default: result = multierror.Append(result, MakeValidationErr(ErrorUnsupportedType, "MessageQueueTriggerSpec.MessageQueueType", spec.MessageQueueType, "not a supported message queue type")) } diff --git a/test/test_utils.sh b/test/test_utils.sh index 8c05e3a6..87bca86d 100755 --- a/test/test_utils.sh +++ b/test/test_utils.sh @@ -435,7 +435,7 @@ dump_logs() { dump_fission_logs $ns $fns executor dump_fission_logs $ns $fns storagesvc dump_fission_logs $ns $fns mqtrigger - dump_fission_logs $ns $fns nats-streaming + dump_fission_logs $ns $fns mqtrigger-nats-streaming dump_function_pod_logs $ns $fns dump_builder_pod_logs $bns dump_fission_crds diff --git a/test/tests/test_package_command.sh b/test/tests/test_package_command.sh index 38c7107b..419093d7 100755 --- a/test/tests/test_package_command.sh +++ b/test/tests/test_package_command.sh @@ -21,6 +21,8 @@ waitBuild() { if [[ $? -eq 0 ]]; then break fi + log "Waiting for build to finish" + sleep 1 done } export -f waitBuild diff --git a/types.go b/types.go index 30efc137..a645a875 100644 --- a/types.go +++ b/types.go @@ -140,8 +140,9 @@ const ( ) const ( - MessageQueueTypeNats = fv1.MessageQueueTypeNats - MessageQueueTypeASQ = fv1.MessageQueueTypeASQ + MessageQueueTypeNats = fv1.MessageQueueTypeNats + MessageQueueTypeASQ = fv1.MessageQueueTypeASQ + MessageQueueTypeKafka = fv1.MessageQueueTypeKafka ) const (