diff --git a/mqtrigger/messageQueue/kafka_test.go b/mqtrigger/messageQueue/kafka_test.go deleted file mode 100644 index 967d4856..00000000 --- a/mqtrigger/messageQueue/kafka_test.go +++ /dev/null @@ -1,119 +0,0 @@ -/* -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/test/tests/mqtrigger/kafka/hellokafka.js b/test/tests/mqtrigger/kafka/hellokafka.js new file mode 100644 index 00000000..22f2dc7d --- /dev/null +++ b/test/tests/mqtrigger/kafka/hellokafka.js @@ -0,0 +1,9 @@ + +module.exports = async function (context) { + console.log(context.request.body); + let obj = context.request.body; + return { + status: 200, + body: obj + }; +} diff --git a/test/tests/mqtrigger/kafka/kafka_pub/glide.lock b/test/tests/mqtrigger/kafka/kafka_pub/glide.lock new file mode 100644 index 00000000..2d3475f6 --- /dev/null +++ b/test/tests/mqtrigger/kafka/kafka_pub/glide.lock @@ -0,0 +1,26 @@ +hash: 1c2fb68841e4a62048d54923d75952578dc4812d9a5c060b031e7c4d47c4e97e +updated: 2018-10-19T14:57:56.260957841+05:30 +imports: +- name: github.com/davecgh/go-spew + version: 8991bc29aa16c548c550c7ff78260e27b9ab7c73 + subpackages: + - spew +- 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/golang/snappy + version: 2e65f85255dbc3072edf28d6b5b8efc472979f5a +- name: github.com/pierrec/lz4 + version: 6b9367c9ff401dbc54fabce3fb8d972e799b702d + subpackages: + - internal/xxh32 +- name: github.com/rcrowley/go-metrics + version: e2704e165165ec55d062f5919b4b29494e9fa790 +- name: github.com/Shopify/sarama + version: a6144ae922fd99dd0ea5046c8137acfb7fab0914 +testImports: [] diff --git a/test/tests/mqtrigger/kafka/kafka_pub/glide.yaml b/test/tests/mqtrigger/kafka/kafka_pub/glide.yaml new file mode 100644 index 00000000..0bb9901b --- /dev/null +++ b/test/tests/mqtrigger/kafka/kafka_pub/glide.yaml @@ -0,0 +1,2 @@ +import: +- package: github.com/Shopify/sarama diff --git a/test/tests/mqtrigger/kafka/kafka_pub/kafka-pub.go b/test/tests/mqtrigger/kafka/kafka_pub/kafka-pub.go new file mode 100644 index 00000000..f9fc565c --- /dev/null +++ b/test/tests/mqtrigger/kafka/kafka_pub/kafka-pub.go @@ -0,0 +1,32 @@ +package main + +import ( + "fmt" + "net/http" + + sarama "github.com/Shopify/sarama" +) + +// Handler posts a message to Kafka Topic +func Handler(w http.ResponseWriter, r *http.Request) { + brokers := []string{"broker.kafka.svc.cluster.local:9092"} + producerConfig := sarama.NewConfig() + producerConfig.Producer.RequiredAcks = sarama.WaitForAll + producerConfig.Producer.Retry.Max = 10 + producerConfig.Producer.Return.Successes = true + producer, err := sarama.NewSyncProducer(brokers, producerConfig) + fmt.Println("Created a new producer ", producer) + if err != nil { + panic(err) + } + _, _, err = producer.SendMessage(&sarama.ProducerMessage{ + Topic: "testtopic", + Value: sarama.StringEncoder("{\"name\": \"testvalue\"}"), + }) + + if err != nil { + w.Write([]byte(fmt.Sprintf("Failed to publish message to topic %s: %v", "testtopic", err))) + return + } + w.Write([]byte("Successfully sent to testtopic")) +} diff --git a/test/tests/mqtrigger/kafka/test_kafka.sh b/test/tests/mqtrigger/kafka/test_kafka.sh new file mode 100755 index 00000000..e96416f7 --- /dev/null +++ b/test/tests/mqtrigger/kafka/test_kafka.sh @@ -0,0 +1,88 @@ +#!/bin/bash +#test:disabled + +# Create a function and trigger it using Kafka +# This test requires Kafka & MQ-Kafka component of Fission installed in the cluster +set -euo pipefail +set +x + +nodeenv="node-kafka" +goenv="go-kafka" +producerfunc="producer-func" +consumerfunc="consumer-func" + +log() { + echo $1 +} +export -f log + +test_mqmessage() { + echo "Checking for valid response" + + while true; do + response0=$(kubectl -nfission logs -l=messagequeue=kafka) + echo $response0 | grep -i $1 + if [[ $? -eq 0 ]]; then + break + fi + sleep 1 + done +} +export -f test_mqmessage + +waitBuild() { + log "Waiting for builder manager to finish the build" + + while true; do + kubectl --namespace default get packages $1 -o jsonpath='{.status.buildstatus}'|grep succeeded + if [[ $? -eq 0 ]]; then + break + fi + log "Waiting for build to finish" + sleep 1 + done +} +export -f waitBuild + +cleanup() { + log "Cleaning up..." + fission env delete --name ${goenv} || true + fission env delete --name ${nodeenv} || true + fission fn delete --name ${producerfunc} || true + fission fn delete --name ${consumerfunc} || true +} +export -f cleanup + +DIR=$(dirname $0) + +log "Creating ${nodeenv} environment" +fission env create --name ${nodeenv} --image fission/node-env +trap cleanup EXIT + +log "Creating ${goenv} environment" +fission env create --name ${goenv} --image fission/go-env --builder fission/go-builder + +log "Creating package for Kafka producer" +pushd $DIR/kafka_pub +zip -qr kafka.zip * +pkgName=$(fission package create --env ${goenv} --src kafka.zip|cut -f2 -d' '| tr -d \') + +log "pkgName=${pkgName}" +popd + +gtimeout 60s bash -c "waitBuild $pkgName" +log "Package ${pkgName} created" + +log "Creating function ${consumerfunc}" +fission fn create --name ${consumerfunc} --env ${nodeenv} --code hellokafka.js + +log "Creating function ${producerfunc}" +fission fn create --name ${producerfunc} --env ${goenv} --pkg ${pkgName} --entrypoint Handler + +log "Creating " +fission mqt create --name kafkatest --function ${consumerfunc} --mqtype kafka --topic testtopic --resptopic resptopic + +fission fn test --name ${producerfunc} + +log "Testing pool manager function" +gtimeout 60 bash -c "test_mqmessage 'testvalue'" \ No newline at end of file diff --git a/test/tests/mqtrigger/main.js b/test/tests/mqtrigger/nats/main.js similarity index 100% rename from test/tests/mqtrigger/main.js rename to test/tests/mqtrigger/nats/main.js diff --git a/test/tests/mqtrigger/main_error.js b/test/tests/mqtrigger/nats/main_error.js similarity index 100% rename from test/tests/mqtrigger/main_error.js rename to test/tests/mqtrigger/nats/main_error.js diff --git a/test/tests/mqtrigger/stan-pub.go b/test/tests/mqtrigger/nats/stan-pub.go similarity index 100% rename from test/tests/mqtrigger/stan-pub.go rename to test/tests/mqtrigger/nats/stan-pub.go diff --git a/test/tests/mqtrigger/stan-sub.go b/test/tests/mqtrigger/nats/stan-sub.go similarity index 100% rename from test/tests/mqtrigger/stan-sub.go rename to test/tests/mqtrigger/nats/stan-sub.go diff --git a/test/tests/mqtrigger/test_mqtrigger.sh b/test/tests/mqtrigger/nats/test_mqtrigger.sh similarity index 98% rename from test/tests/mqtrigger/test_mqtrigger.sh rename to test/tests/mqtrigger/nats/test_mqtrigger.sh index 5f5113e7..f9054a65 100755 --- a/test/tests/mqtrigger/test_mqtrigger.sh +++ b/test/tests/mqtrigger/nats/test_mqtrigger.sh @@ -7,7 +7,7 @@ set -euo pipefail set +x -ROOT=$(dirname $0)/../.. +ROOT=$(dirname $0)/../../.. DIR=$(dirname $0) clusterID="fissionMQTrigger" diff --git a/test/tests/mqtrigger/test_mqtrigger_error.sh b/test/tests/mqtrigger/nats/test_mqtrigger_error.sh similarity index 98% rename from test/tests/mqtrigger/test_mqtrigger_error.sh rename to test/tests/mqtrigger/nats/test_mqtrigger_error.sh index e4274498..1e062556 100755 --- a/test/tests/mqtrigger/test_mqtrigger_error.sh +++ b/test/tests/mqtrigger/nats/test_mqtrigger_error.sh @@ -7,7 +7,7 @@ set -euo pipefail set +x -ROOT=$(dirname $0)/../.. +ROOT=$(dirname $0)/../../.. DIR=$(dirname $0) clusterID="fissionMQTrigger"