Added support for Kafka headers to be passed and retrieved from functions. The headers and supported and work only for Kafka version 0.11.0.0 and higher.
236 lines
6.8 KiB
Go
236 lines
6.8 KiB
Go
/*
|
|
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"
|
|
"os"
|
|
"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
|
|
version sarama.KafkaVersion
|
|
}
|
|
)
|
|
|
|
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")
|
|
}
|
|
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 %v", 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 %q", trigger)
|
|
log.Infof("brokers set to %q", kafka.brokers)
|
|
|
|
// Create new consumer
|
|
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 {
|
|
panic(err)
|
|
}
|
|
|
|
// Create new producer
|
|
producerConfig := sarama.NewConfig()
|
|
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 {
|
|
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, msg) {
|
|
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, 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",
|
|
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)
|
|
|
|
// 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 {
|
|
log.Warningf("Request creation failed: %v", url)
|
|
return false
|
|
}
|
|
|
|
// 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++ {
|
|
// Make the request
|
|
resp, err = http.DefaultClient.Do(req)
|
|
if err != nil {
|
|
log.Errorf("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 {
|
|
// 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),
|
|
Headers: kafkaRecordHeaders,
|
|
})
|
|
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)
|
|
}
|
|
}
|