diff --git a/charts/README.md b/charts/README.md index 2d9b2d6f..a73d9192 100644 --- a/charts/README.md +++ b/charts/README.md @@ -51,19 +51,20 @@ The following table lists the configurable parameters of the Fission chart and t * Extra configuration for `fission-all` -| Parameter | Description | Default | -| ------------------------------- | --------------------------- | ---------------------------------------------------------- | -| `logger.influxdbAdmin` | Log database admin username | `admin` | -| `logger.fluentdImage` | Logger fluentd image | `fission/fluentd` | -| `fissionUiImage` | Fission ui image | `fission/fission-ui:0.1.0` | -| `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`) | +| Parameter | Description | Default | +| ------------------------------- | --------------------------- | ---------------------------------------------------------- | +| `logger.influxdbAdmin` | Log database admin username | `admin` | +| `logger.fluentdImage` | Logger fluentd image | `fission/fluentd` | +| `fissionUiImage` | Fission ui image | `fission/fission-ui:0.1.0` | +| `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`) | +| `kafka.version` | Kafka broker version | `None` (should be `>= 0.11.0.0` to enable Kafka record headers support) | * - Please note that deploying of Azure Storage Queue or Kafka is not done by Fission chart and you will have to explicitly deploy them. diff --git a/charts/fission-all/templates/deployment.yaml b/charts/fission-all/templates/deployment.yaml index f18417f8..8c5f3970 100644 --- a/charts/fission-all/templates/deployment.yaml +++ b/charts/fission-all/templates/deployment.yaml @@ -578,6 +578,8 @@ spec: value: kafka - name: MESSAGE_QUEUE_URL value: "{{.Values.kafka.brokers}}" + - name: MESSAGE_QUEUE_KAFKA_VERSION + value: "{{.Values.kafka.version}}" serviceAccount: fission-svc {{- end }} --- diff --git a/charts/fission-all/values.yaml b/charts/fission-all/values.yaml index 273982a2..7c5d2d37 100644 --- a/charts/fission-all/values.yaml +++ b/charts/fission-all/values.yaml @@ -73,6 +73,10 @@ azureStorageQueue: kafka: enabled: false brokers: 'broker.kafka:9092' + ## version of Kafka broker + ## Must be a string in the format "major.minor.veryMinor.patch" + ## Should be >= 0.11.0.0 to enable Kafka record headers support + # version: "0.11.0.0" ## Persist data to a persistent volume. persistence: diff --git a/mqtrigger/messageQueue/kafka.go b/mqtrigger/messageQueue/kafka.go index 018a589f..3f657265 100644 --- a/mqtrigger/messageQueue/kafka.go +++ b/mqtrigger/messageQueue/kafka.go @@ -21,6 +21,7 @@ import ( "fmt" "io/ioutil" "net/http" + "os" "strings" sarama "github.com/Shopify/sarama" @@ -35,6 +36,7 @@ type ( Kafka struct { routerUrl string brokers []string + version sarama.KafkaVersion } ) @@ -42,11 +44,20 @@ func makeKafkaMessageQueue(routerUrl string, mqCfg MessageQueueConfig) (MessageQ if len(routerUrl) == 0 || len(mqCfg.Url) == 0 { return nil, errors.New("The router URL or MQ URL is empty") } + mqKafkaVersion := os.Getenv("MESSAGE_QUEUE_KAFKA_VERSION") + + // Parse version string + kafkaVersion, err := sarama.ParseKafkaVersion(mqKafkaVersion) + if err != nil { + log.Warningf("Error parsing version string %q: %v. Falling back to %q", mqKafkaVersion, err, kafkaVersion) + } + kafka := Kafka{ routerUrl: routerUrl, brokers: strings.Split(mqCfg.Url, ","), + version: kafkaVersion, } - log.Infof("Created Queue %q", kafka) + log.Infof("Created Queue %v", kafka) return kafka, nil } @@ -62,6 +73,7 @@ func (kafka Kafka) subscribe(trigger *crd.MessageQueueTrigger) (messageQueueSubs consumerConfig := cluster.NewConfig() consumerConfig.Consumer.Return.Errors = true consumerConfig.Group.Return.Notifications = true + consumerConfig.Config.Version = kafka.version consumer, err := cluster.NewConsumer(kafka.brokers, string(trigger.Metadata.UID), []string{trigger.Spec.Topic}, consumerConfig) log.Infof("Created a new consumer: %#v", consumer) if err != nil { @@ -73,6 +85,7 @@ func (kafka Kafka) subscribe(trigger *crd.MessageQueueTrigger) (messageQueueSubs producerConfig.Producer.RequiredAcks = sarama.WaitForAll producerConfig.Producer.Retry.Max = 10 producerConfig.Producer.Return.Successes = true + producerConfig.Version = kafka.version producer, err := sarama.NewSyncProducer(kafka.brokers, producerConfig) log.Infof("Created a new producer %q", producer) if err != nil { @@ -97,7 +110,7 @@ func (kafka Kafka) subscribe(trigger *crd.MessageQueueTrigger) (messageQueueSubs 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[:])) { + if kafkaMsgHandler(&kafka, producer, trigger, msg) { consumer.MarkOffset(msg, "") // mark message as processed } } @@ -110,7 +123,8 @@ 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 { +func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *crd.MessageQueueTrigger, msg *sarama.ConsumerMessage) bool { + var value string = string(msg.Value[:]) // Support other function ref types if trigger.Spec.FunctionReference.Type != fission.FunctionReferenceTypeFunctionName { log.Fatalf("Unsupported function reference type (%v) for trigger %v", @@ -119,12 +133,15 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *crd.Me 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{ + + // Generate the Headers + fissionHeaders := 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 { @@ -132,9 +149,16 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *crd.Me return false } - for k, v := range headers { + // Set the headers came from Kafka record + // Using Header.Add() as msg.Headers may have keys with more than one value + for _, h := range msg.Headers { + req.Header.Add(string(h.Key), string(h.Value)) + } + + for k, v := range fissionHeaders { req.Header.Set(k, v) } + // Make the request var resp *http.Response for attempt := 0; attempt <= trigger.Spec.MaxRetries; attempt++ { @@ -169,9 +193,23 @@ func kafkaMsgHandler(kafka *Kafka, producer sarama.SyncProducer, trigger *crd.Me return false } if len(trigger.Spec.ResponseTopic) > 0 { + // Generate Kafka record headers + var kafkaRecordHeaders []sarama.RecordHeader + if kafka.version.IsAtLeast(sarama.V0_11_0_0) { + for k, v := range resp.Header { + // One key may have multiple values + for _, v := range v { + kafkaRecordHeaders = append(kafkaRecordHeaders, sarama.RecordHeader{Key: []byte(k), Value: []byte(v)}) + } + } + } else { + log.Warningf("Headers are not supported by Kafka version %q, needs v0.11+: dropping the headers", kafka.version) + } + _, _, err := producer.SendMessage(&sarama.ProducerMessage{ - Topic: trigger.Spec.ResponseTopic, - Value: sarama.StringEncoder(body), + Topic: trigger.Spec.ResponseTopic, + Value: sarama.StringEncoder(body), + Headers: kafkaRecordHeaders, }) if err != nil { log.Warningf("Failed to publish message to topic %s: %v", trigger.Spec.ResponseTopic, err) diff --git a/test/tests/mqtrigger/kafka/hellokafka.js b/test/tests/mqtrigger/kafka/hellokafka.js index 22f2dc7d..d9b64ac6 100644 --- a/test/tests/mqtrigger/kafka/hellokafka.js +++ b/test/tests/mqtrigger/kafka/hellokafka.js @@ -1,9 +1,13 @@ module.exports = async function (context) { console.log(context.request.body); + console.log("z-custom-name: " + context.request.headers['z-custom-name']); + console.log("x-fission-function-name: " + context.request.headers['x-fission-function-name']); let obj = context.request.body; + let headers = context.request.headers; return { status: 200, + headers: headers, body: obj }; } diff --git a/test/tests/mqtrigger/kafka/kafka_pub/kafka-pub.go b/test/tests/mqtrigger/kafka/kafka_pub/kafka-pub.go index f9fc565c..c04cc78e 100644 --- a/test/tests/mqtrigger/kafka/kafka_pub/kafka-pub.go +++ b/test/tests/mqtrigger/kafka/kafka_pub/kafka-pub.go @@ -14,14 +14,18 @@ func Handler(w http.ResponseWriter, r *http.Request) { producerConfig.Producer.RequiredAcks = sarama.WaitForAll producerConfig.Producer.Retry.Max = 10 producerConfig.Producer.Return.Successes = true + producerConfig.Version = sarama.V0_11_0_0 producer, err := sarama.NewSyncProducer(brokers, producerConfig) fmt.Println("Created a new producer ", producer) if err != nil { panic(err) } + + headers := []sarama.RecordHeader{{Key: []byte("Z-Custom-Name"), Value: []byte("Kafka-Header-test")}} _, _, err = producer.SendMessage(&sarama.ProducerMessage{ - Topic: "testtopic", - Value: sarama.StringEncoder("{\"name\": \"testvalue\"}"), + Topic: "testtopic", + Value: sarama.StringEncoder("{\"name\": \"testvalue\"}"), + Headers: headers, }) if err != nil { diff --git a/test/tests/mqtrigger/kafka/test_kafka.sh b/test/tests/mqtrigger/kafka/test_kafka.sh index e96416f7..0a55e39c 100755 --- a/test/tests/mqtrigger/kafka/test_kafka.sh +++ b/test/tests/mqtrigger/kafka/test_kafka.sh @@ -10,6 +10,7 @@ nodeenv="node-kafka" goenv="go-kafka" producerfunc="producer-func" consumerfunc="consumer-func" +consumerfunc2="consumer-func2" log() { echo $1 @@ -30,6 +31,23 @@ test_mqmessage() { } export -f test_mqmessage +test_fnmessage() { + # $1: functionName + # $2: container name + # $3: string to look for + echo "Checking for valid function log" + + while true; do + response0=$(kubectl -nfission-function logs -l=functionName=$1 -c $2) + echo $response0 | grep -i "$3" + if [[ $? -eq 0 ]]; then + break + fi + sleep 1 + done +} +export -f test_fnmessage + waitBuild() { log "Waiting for builder manager to finish the build" @@ -50,6 +68,9 @@ cleanup() { fission env delete --name ${nodeenv} || true fission fn delete --name ${producerfunc} || true fission fn delete --name ${consumerfunc} || true + fission fn delete --name ${consumerfunc2} || true + fission mqt delete --name kafkatest || true + fission mqt delete --name kafkatest2 || true } export -f cleanup @@ -64,25 +85,40 @@ fission env create --name ${goenv} --image fission/go-env --builder fission/go-b log "Creating package for Kafka producer" pushd $DIR/kafka_pub +glide install 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" +timeout 120s 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 ${consumerfunc2}" +fission fn create --name ${consumerfunc2} --env ${nodeenv} --code hellokafka.js + log "Creating function ${producerfunc}" fission fn create --name ${producerfunc} --env ${goenv} --pkg ${pkgName} --entrypoint Handler -log "Creating " +log "Creating trigger kafkatest" fission mqt create --name kafkatest --function ${consumerfunc} --mqtype kafka --topic testtopic --resptopic resptopic +log "Creating trigger kafkatest2" +fission mqt create --name kafkatest2 --function ${consumerfunc2} --mqtype kafka --topic resptopic + fission fn test --name ${producerfunc} log "Testing pool manager function" -gtimeout 60 bash -c "test_mqmessage 'testvalue'" \ No newline at end of file +timeout 60 bash -c "test_mqmessage 'testvalue'" + +log "Testing the headers values in ${consumerfunc}" +timeout 60 bash -c "test_fnmessage '${consumerfunc}' '${nodeenv}' 'z-custom-name: Kafka-Header-test'" + +log "Testing the header value in ${consumerfunc2}" +timeout 60 bash -c "test_fnmessage '${consumerfunc2}' '${nodeenv}' 'z-custom-name: Kafka-Header-test'" +# test if the Fission specific headers are overwritten +timeout 60 bash -c "test_fnmessage '${consumerfunc2}' '${nodeenv}' 'x-fission-function-name: consumer-func2'"