From 395a8adf371fcc3c212c92c554a40f997bc499bb Mon Sep 17 00:00:00 2001 From: Suraj Banakar <34534103+vadasambar@users.noreply.github.com> Date: Wed, 9 Oct 2019 21:48:57 +0530 Subject: [PATCH] Implement TLS authentication for kafka mqt (#1300) * use secrets to store keys and certificates --- charts/fission-all/templates/deployment.yaml | 43 +++++++++++ charts/fission-all/values.yaml | 20 +++++- pkg/mqtrigger/messageQueue/kafka.go | 75 +++++++++++++++++--- pkg/mqtrigger/messageQueue/messageQueue.go | 5 +- pkg/mqtrigger/mqtrigger.go | 56 ++++++++++++++- 5 files changed, 184 insertions(+), 15 deletions(-) diff --git a/charts/fission-all/templates/deployment.yaml b/charts/fission-all/templates/deployment.yaml index 0c0cded2..8eba3f11 100644 --- a/charts/fission-all/templates/deployment.yaml +++ b/charts/fission-all/templates/deployment.yaml @@ -632,7 +632,50 @@ spec: value: {{ .Values.traceSamplingRate | default "0.5" | quote }} - name: DEBUG_ENV value: {{ .Values.debugEnv | quote }} + # TLS authentication is TLS with authentication (2 way) + # More info: https://docs.confluent.io/current/kafka/authentication_ssl.html#ssl-overview + {{- if .Values.kafka.authentication.tls.enabled }} + - name: TLS_ENABLED + value: "true" + - name: MESSAGE_QUEUE_SECRETS + value: /etc/fission/secrets + volumeMounts: + - name: kafka-secrets + mountPath: /etc/fission/secrets + {{- end }} serviceAccount: fission-svc + {{- if .Values.kafka.authentication.tls.enabled }} + volumes: + - name: kafka-secrets + secret: + secretName: mqtrigger-kafka-secrets + {{- end }} + +--- +{{- if .Values.kafka.authentication.tls.enabled }} +apiVersion: v1 +kind: Secret +metadata: + name: mqtrigger-kafka-secrets + labels: + chart: "{{ .Chart.Name }}-{{ .Chart.Version }}" +data: + {{- if .Files.Get (printf "%s" .Values.kafka.authentication.tls.caCert) }} + caCert: {{ .Files.Get (printf "%s" .Values.kafka.authentication.tls.caCert) | b64enc }} + {{- else }} + {{ fail "Invalid chart. CA Certificate not found." }} + {{- end }} + {{- if .Files.Get (printf "%s" .Values.kafka.authentication.tls.userCert) }} + userCert: {{ .Files.Get (printf "%s" .Values.kafka.authentication.tls.userCert) | b64enc }} + {{- else }} + {{ fail "Invalid chart. User Certificate not found." }} + {{- end }} + {{- if .Files.Get (printf "%s" .Values.kafka.authentication.tls.userKey) }} + userKey: {{ .Files.Get (printf "%s" .Values.kafka.authentication.tls.userKey) | b64enc }} + {{- else }} + {{ fail "Invalid chart. User Key not found." }} + {{- end }} +{{- end }} {{- if .Values.extraCoreComponentPodConfig }} {{ toYaml .Values.extraCoreComponentPodConfig | indent 6 -}} {{- end }} diff --git a/charts/fission-all/values.yaml b/charts/fission-all/values.yaml index 2ecc4665..0bc53c25 100644 --- a/charts/fission-all/values.yaml +++ b/charts/fission-all/values.yaml @@ -108,7 +108,25 @@ azureStorageQueue: ## Kafka: enable and configure the details kafka: enabled: false - brokers: 'broker.kafka:9092' + # note: below link is only for reference. + # Please use the brokers link for your kafka here. + brokers: 'broker.kafka:9092' # or your-bootstrap-server.kafka:9092/9093 + authentication: + tls: + enabled: false + caCert: '' # path to certificate containing public key of CA authority + userCert: '' # path to certificate containing public key of the user signed by CA authority + userKey: '' # path to private key of the user + + # brokers: 'my-broker.kafka:9092' # or my-bootstrap-server.kafka:9092/9093 + # Sample config for authentication + # authentication: + # tls: + # enabled: true + # caCert: 'auth/kafka/ca.crt' + # userCert: 'auth/kafka/user.crt' + # userKey: 'auth/kafka/user.key' + ## version of Kafka broker ## For 0.x it must be a string in the format ## "major.minor.veryMinor.patch" example: 0.8.2.0 diff --git a/pkg/mqtrigger/messageQueue/kafka.go b/pkg/mqtrigger/messageQueue/kafka.go index dac39eff..4d41cabe 100644 --- a/pkg/mqtrigger/messageQueue/kafka.go +++ b/pkg/mqtrigger/messageQueue/kafka.go @@ -17,10 +17,13 @@ limitations under the License. package messageQueue import ( + "crypto/tls" + "crypto/x509" "fmt" "io/ioutil" "net/http" "os" + "strconv" "strings" sarama "github.com/Shopify/sarama" @@ -39,6 +42,8 @@ type ( routerUrl string brokers []string version sarama.KafkaVersion + authKeys map[string][]byte + tls bool } ) @@ -64,9 +69,23 @@ func makeKafkaMessageQueue(logger *zap.Logger, routerUrl string, mqCfg MessageQu version: kafkaVersion, } + if tls, _ := strconv.ParseBool(os.Getenv("TLS_ENABLED")); tls == true { + kafka.tls = true + + authKeys := make(map[string][]byte) + + if mqCfg.Secrets == nil { + return nil, errors.New("no secrets were loaded") + } + + authKeys["caCert"] = mqCfg.Secrets["caCert"] + authKeys["userCert"] = mqCfg.Secrets["userCert"] + authKeys["userKey"] = mqCfg.Secrets["userKey"] + kafka.authKeys = authKeys + } + logger.Info("created kafka queue", zap.Any("kafka brokers", kafka.brokers), zap.Any("kafka version", kafka.version)) - return kafka, nil } @@ -83,6 +102,28 @@ func (kafka Kafka) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubs consumerConfig.Consumer.Return.Errors = true consumerConfig.Group.Return.Notifications = true consumerConfig.Config.Version = kafka.version + + // 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 + + // Setup TLS for both producer and consumer + if kafka.tls { + consumerConfig.Net.TLS.Enable = true + producerConfig.Net.TLS.Enable = true + tlsConfig, err := kafka.getTLSConfig() + + if err != nil { + return nil, err + } + + producerConfig.Net.TLS.Config = tlsConfig + consumerConfig.Net.TLS.Config = tlsConfig + } + consumer, err := cluster.NewConsumer(kafka.brokers, string(trigger.Metadata.UID), []string{trigger.Spec.Topic}, consumerConfig) kafka.logger.Info("created a new consumer", zap.Strings("brokers", kafka.brokers), zap.String("input topic", trigger.Spec.Topic), @@ -91,17 +132,10 @@ func (kafka Kafka) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubs zap.String("trigger name", trigger.Metadata.Name), zap.String("function namespace", trigger.Metadata.Namespace), zap.String("function name", trigger.Spec.FunctionReference.Name)) - if err != nil { - panic(err) + return nil, 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) kafka.logger.Info("created a new producer", zap.Strings("brokers", kafka.brokers), zap.String("input topic", trigger.Spec.Topic), @@ -112,7 +146,7 @@ func (kafka Kafka) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubs zap.String("function name", trigger.Spec.FunctionReference.Name)) if err != nil { - panic(err) + return nil, err } // consume errors @@ -142,6 +176,27 @@ func (kafka Kafka) subscribe(trigger *fv1.MessageQueueTrigger) (messageQueueSubs return consumer, nil } +func (kafka Kafka) getTLSConfig() (*tls.Config, error) { + tlsConfig := tls.Config{} + cert, err := tls.X509KeyPair(kafka.authKeys["userCert"], kafka.authKeys["userKey"]) + if err != nil { + return nil, err + } + + tlsConfig.Certificates = []tls.Certificate{cert} + + if err != nil { + return nil, err + } + + caCertPool := x509.NewCertPool() + caCertPool.AppendCertsFromPEM(kafka.authKeys["caCert"]) + tlsConfig.RootCAs = caCertPool + tlsConfig.BuildNameToCertificate() + + return &tlsConfig, nil +} + func (kafka Kafka) unsubscribe(subscription messageQueueSubscription) error { return subscription.(*cluster.Consumer).Close() } diff --git a/pkg/mqtrigger/messageQueue/messageQueue.go b/pkg/mqtrigger/messageQueue/messageQueue.go index a436ad09..9f9597c4 100644 --- a/pkg/mqtrigger/messageQueue/messageQueue.go +++ b/pkg/mqtrigger/messageQueue/messageQueue.go @@ -42,8 +42,9 @@ type ( requestType int MessageQueueConfig struct { - MQType string - Url string + MQType string + Url string + Secrets map[string][]byte } MessageQueue interface { diff --git a/pkg/mqtrigger/mqtrigger.go b/pkg/mqtrigger/mqtrigger.go index dba5faa7..9012cf15 100644 --- a/pkg/mqtrigger/mqtrigger.go +++ b/pkg/mqtrigger/mqtrigger.go @@ -17,7 +17,11 @@ limitations under the License. package mqtrigger import ( + "fmt" + "io/ioutil" "os" + "path" + "strings" "github.com/pkg/errors" "go.uber.org/zap" @@ -28,6 +32,7 @@ import ( func Start(logger *zap.Logger, routerUrl string) error { fissionClient, _, _, err := crd.MakeFissionClient() + if err != nil { return errors.Wrap(err, "failed to get fission or kubernetes client") } @@ -40,10 +45,57 @@ func Start(logger *zap.Logger, routerUrl string) error { // Message queue type: nats is the only supported one for now mqType := os.Getenv("MESSAGE_QUEUE_TYPE") mqUrl := os.Getenv("MESSAGE_QUEUE_URL") + + secretsPath := strings.TrimSpace(os.Getenv("MESSAGE_QUEUE_SECRETS")) + + var secrets map[string][]byte + if len(secretsPath) > 0 { + // For authentication with message queue + secrets, err = readSecrets(logger, secretsPath) + if err != nil { + return err + } + } + mqCfg := messageQueue.MessageQueueConfig{ - MQType: mqType, - Url: mqUrl, + MQType: mqType, + Url: mqUrl, + Secrets: secrets, } messageQueue.MakeMessageQueueTriggerManager(logger, fissionClient, routerUrl, mqCfg) return nil } + +func readSecrets(logger *zap.Logger, secretsPath string) (map[string][]byte, error) { + + // return if no secrets exist + if _, err := os.Stat(secretsPath); os.IsNotExist(err) { + return nil, err + } + + secretFiles, err := ioutil.ReadDir(secretsPath) + if err != nil { + return nil, err + } + + secrets := make(map[string][]byte) + for _, secretFile := range secretFiles { + + fileName := secretFile.Name() + // /etc/secrets contain some hidden directories (like .data) + // ignore them + if !secretFile.IsDir() && !strings.HasPrefix(fileName, ".") { + logger.Info(fmt.Sprintf("Reading secret from %s", fileName)) + + filePath := path.Join(secretsPath, fileName) + secret, fileReadErr := ioutil.ReadFile(filePath) + if fileReadErr != nil { + return nil, fileReadErr + } + + secrets[fileName] = secret + } + } + + return secrets, nil +}