Fission MQT integration with keda (#1657)

Fission MQT integration based on Keda with Kafka connector
This commit is contained in:
Rahul Bhati
2020-07-14 21:35:48 +05:30
committed by GitHub
parent d6fc9f71d0
commit 8e0571c7c9
25 changed files with 2203 additions and 28 deletions
+28
View File
@@ -645,6 +645,34 @@ type (
// Content type of payload
ContentType string `json:"contentType"`
// The period to check each trigger source on every ScaledObject, and scale the deployment up or down accordingly
// +optional
PollingInterval *int32 `json:"pollingInterval,omitempty"`
// The period to wait after the last trigger reported active before scaling the deployment back to 0
// +optional
CooldownPeriod *int32 `json:"cooldownPeriod,omitempty"`
// Minimum number of replicas KEDA will scale the deployment down to
// +optional
MinReplicaCount *int32 `json:"minReplicaCount,omitempty"`
// Maximum number of replicas KEDA will scale the deployment up to
// +optional
MaxReplicaCount *int32 `json:"maxReplicaCount,omitempty"`
// ScalerTrigger fields
// +optional
Metadata map[string]string `json:"metadata"`
// Secret name
// +optional
Secret string `json:"secret,omitempty"`
// Kind of Message Queue Trigger to be created, by default its fission
// +optional
MqtKind string `json:"mqtkind,omitempty"`
}
// TimeTrigger invokes the specific function at a time or
+27
View File
@@ -710,6 +710,33 @@ func (in *MessageQueueTriggerList) DeepCopyObject() runtime.Object {
func (in *MessageQueueTriggerSpec) DeepCopyInto(out *MessageQueueTriggerSpec) {
*out = *in
in.FunctionReference.DeepCopyInto(&out.FunctionReference)
if in.PollingInterval != nil {
in, out := &in.PollingInterval, &out.PollingInterval
*out = new(int32)
**out = **in
}
if in.CooldownPeriod != nil {
in, out := &in.CooldownPeriod, &out.CooldownPeriod
*out = new(int32)
**out = **in
}
if in.MinReplicaCount != nil {
in, out := &in.MinReplicaCount, &out.MinReplicaCount
*out = new(int32)
**out = **in
}
if in.MaxReplicaCount != nil {
in, out := &in.MaxReplicaCount, &out.MaxReplicaCount
*out = new(int32)
**out = **in
}
if in.Metadata != nil {
in, out := &in.Metadata, &out.Metadata
*out = make(map[string]string, len(*in))
for key, val := range *in {
(*out)[key] = val
}
}
return
}
+27
View File
@@ -23,6 +23,7 @@ import (
apiextensionsclient "k8s.io/apiextensions-apiserver/pkg/client/clientset/clientset"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/kubernetes"
_ "k8s.io/client-go/plugin/pkg/client/auth"
"k8s.io/client-go/rest"
@@ -107,3 +108,29 @@ func (fc *FissionClient) WaitForCRDs() error {
}
}
}
// GetDynamicClient creates and returns new dynamic client or returns an error
func GetDynamicClient() (dynamic.Interface, error) {
var config *rest.Config
var err error
// get the config, either from kubeconfig or using our
// in-cluster service account
kubeConfig := os.Getenv("KUBECONFIG")
if len(kubeConfig) != 0 {
config, err = clientcmd.BuildConfigFromFlags("", kubeConfig)
if err != nil {
return nil, err
}
} else {
config, err = rest.InClusterConfig()
if err != nil {
return nil, err
}
}
dynamicClient, err := dynamic.NewForConfig(config)
if err != nil {
return nil, err
}
return dynamicClient, nil
}
+6 -2
View File
@@ -33,7 +33,9 @@ func Commands() *cobra.Command {
Required: []flag.Flag{flag.MqtFnName, flag.MqtTopic},
Optional: []flag.Flag{flag.MqtName, flag.MqtMQType, flag.MqtRespTopic,
flag.MqtErrorTopic, flag.MqtMaxRetries, flag.MqtMsgContentType,
flag.NamespaceFunction, flag.SpecSave, flag.SpecDry},
flag.NamespaceFunction, flag.SpecSave, flag.SpecDry, flag.MqtPollingInterval,
flag.MqtCooldownPeriod, flag.MqtMinReplicaCount, flag.MqtMaxReplicaCount, flag.MqtSecret,
flag.MqtMetadata, flag.MqtKind},
})
updateCmd := &cobra.Command{
@@ -45,7 +47,9 @@ func Commands() *cobra.Command {
wrapper.SetFlags(updateCmd, flag.FlagSet{
Required: []flag.Flag{flag.MqtName},
Optional: []flag.Flag{flag.MqtFnName, flag.MqtTopic, flag.MqtRespTopic, flag.MqtErrorTopic,
flag.MqtMaxRetries, flag.MqtMsgContentType, flag.NamespaceTrigger},
flag.MqtMaxRetries, flag.MqtMsgContentType, flag.NamespaceTrigger, flag.MqtPollingInterval,
flag.MqtCooldownPeriod, flag.MqtMinReplicaCount, flag.MqtMaxReplicaCount, flag.MqtMetadata,
flag.MqtSecret, flag.MqtKind},
})
deleteCmd := &cobra.Command{
+35
View File
@@ -93,6 +93,34 @@ func (opts *CreateSubCommand) complete(input cli.Input) error {
return err
}
pollingInterval := int32(input.Int(flagkey.MqtPollingInterval))
if pollingInterval < 0 {
return errors.New("Polling interval must be greater than or equal to 0")
}
cooldownPeriod := int32(input.Int(flagkey.MqtCooldownPeriod))
if cooldownPeriod < 0 {
return errors.New("CooldownPeriod interval is the period to wait after the last trigger reported active before scaling the deployment back to 0, it must be greater than or equal to 0")
}
minReplicaCount := int32(input.Int(flagkey.MqtMinReplicaCount))
if minReplicaCount < 0 {
return errors.New("MinReplicaCount must be greater than or equal to 0")
}
maxReplicaCount := int32(input.Int(flagkey.MqtMaxReplicaCount))
if maxReplicaCount < 0 {
return errors.New("MaxReplicaCount must be greater than or equal to 0")
}
metadata := make(map[string]string)
metadataParams := input.StringSlice(flagkey.MqtMetadata)
_ = util.UpdateMapFromStringSlice(&metadata, metadataParams)
secret := input.String(flagkey.MqtSecret)
mqtKind := input.String(flagkey.MqtKind)
if input.Bool(flagkey.SpecSave) {
specDir := util.GetSpecDir(input)
fr, err := spec.ReadSpecs(specDir)
@@ -131,6 +159,13 @@ func (opts *CreateSubCommand) complete(input cli.Input) error {
ErrorTopic: errorTopic,
MaxRetries: maxRetries,
ContentType: contentType,
PollingInterval: &pollingInterval,
CooldownPeriod: &cooldownPeriod,
MinReplicaCount: &minReplicaCount,
MaxReplicaCount: &maxReplicaCount,
Metadata: metadata,
Secret: secret,
MqtKind: mqtKind,
},
}
+39 -3
View File
@@ -26,6 +26,7 @@ import (
"github.com/fission/fission/pkg/fission-cli/cliwrapper/cli"
"github.com/fission/fission/pkg/fission-cli/cmd"
flagkey "github.com/fission/fission/pkg/fission-cli/flag/key"
"github.com/fission/fission/pkg/fission-cli/util"
)
type UpdateSubCommand struct {
@@ -60,7 +61,13 @@ func (opts *UpdateSubCommand) complete(input cli.Input) error {
maxRetries := input.Int(flagkey.MqtMaxRetries)
fnName := input.String(flagkey.MqtFnName)
contentType := input.String(flagkey.MqtMsgContentType)
pollingInterval := int32(input.Int(flagkey.MqtPollingInterval))
cooldownPeriod := int32(input.Int(flagkey.MqtCooldownPeriod))
minReplicaCount := int32(input.Int(flagkey.MqtMinReplicaCount))
maxReplicaCount := int32(input.Int(flagkey.MqtMaxReplicaCount))
metadataParams := input.StringSlice(flagkey.MqtMetadata)
secret := input.String(flagkey.MqtSecret)
mqtKind := input.String(flagkey.MqtKind)
// TODO : Find out if we can make a call to checkIfFunctionExists, in the same ns more importantly.
err = checkMQTopicAvailability(mqt.Spec.MessageQueueType, topic, respTopic)
@@ -81,7 +88,7 @@ func (opts *UpdateSubCommand) complete(input cli.Input) error {
mqt.Spec.ErrorTopic = errorTopic
updated = true
}
if maxRetries > -1 {
if input.IsSet(flagkey.MqtMaxRetries) {
mqt.Spec.MaxRetries = maxRetries
updated = true
}
@@ -89,10 +96,39 @@ func (opts *UpdateSubCommand) complete(input cli.Input) error {
mqt.Spec.FunctionReference.Name = fnName
updated = true
}
if len(contentType) > 0 {
if input.IsSet(flagkey.MqtMsgContentType) {
mqt.Spec.ContentType = contentType
updated = true
}
if input.IsSet(flagkey.MqtPollingInterval) {
mqt.Spec.PollingInterval = &pollingInterval
updated = true
}
if input.IsSet(flagkey.MqtCooldownPeriod) {
mqt.Spec.CooldownPeriod = &cooldownPeriod
updated = true
}
if input.IsSet(flagkey.MqtMinReplicaCount) {
mqt.Spec.MinReplicaCount = &minReplicaCount
updated = true
}
if input.IsSet(flagkey.MqtMaxReplicaCount) {
mqt.Spec.MaxReplicaCount = &maxReplicaCount
updated = true
}
if input.IsSet(flagkey.MqtMetadata) {
updated = updated || util.UpdateMapFromStringSlice(&mqt.Spec.Metadata, metadataParams)
}
if input.IsSet(flagkey.MqtSecret) {
mqt.Spec.Secret = secret
updated = true
}
if input.IsSet(flagkey.MqtKind) {
mqt.Spec.MqtKind = mqtKind
updated = true
}
if !updated {
return errors.New("Nothing changed, see 'help' for more details")
+15 -8
View File
@@ -128,14 +128,21 @@ var (
TtFnName = Flag{Type: String, Name: flagkey.TtFnName, Usage: "Function name"}
TtRound = Flag{Type: Int, Name: flagkey.TtRound, Usage: "Get next N rounds of invocation time", DefaultValue: 1}
MqtName = Flag{Type: String, Name: flagkey.MqtName, Usage: "Message queue trigger name"}
MqtFnName = Flag{Type: String, Name: flagkey.MqtFnName, Usage: "Function name"}
MqtMQType = Flag{Type: String, Name: flagkey.MqtMQType, Usage: "Message queue type, e.g. nats-streaming, azure-storage-queue, kafka", DefaultValue: "nats-streaming"}
MqtTopic = Flag{Type: String, Name: flagkey.MqtTopic, Usage: "Message queue Topic the trigger listens on"}
MqtRespTopic = Flag{Type: String, Name: flagkey.MqtRespTopic, Usage: "Topic that the function response is sent on (response discarded if unspecified)"}
MqtErrorTopic = Flag{Type: String, Name: flagkey.MqtErrorTopic, Usage: "Topic that the function error messages are sent to (errors discarded if unspecified"}
MqtMaxRetries = Flag{Type: Int, Name: flagkey.MqtMaxRetries, Usage: "Maximum number of times the function will be retried upon failure", DefaultValue: 0}
MqtMsgContentType = Flag{Type: String, Name: flagkey.MqtMsgContentType, Short: "c", Usage: "Content type of messages that publish to the topic", DefaultValue: "application/json"}
MqtName = Flag{Type: String, Name: flagkey.MqtName, Usage: "Message queue trigger name"}
MqtFnName = Flag{Type: String, Name: flagkey.MqtFnName, Usage: "Function name"}
MqtMQType = Flag{Type: String, Name: flagkey.MqtMQType, Usage: "Message queue type, e.g. nats-streaming, azure-storage-queue, kafka", DefaultValue: "nats-streaming"}
MqtTopic = Flag{Type: String, Name: flagkey.MqtTopic, Usage: "Message queue Topic the trigger listens on"}
MqtRespTopic = Flag{Type: String, Name: flagkey.MqtRespTopic, Usage: "Topic that the function response is sent on (response discarded if unspecified)"}
MqtErrorTopic = Flag{Type: String, Name: flagkey.MqtErrorTopic, Usage: "Topic that the function error messages are sent to (errors discarded if unspecified"}
MqtMaxRetries = Flag{Type: Int, Name: flagkey.MqtMaxRetries, Usage: "Maximum number of times the function will be retried upon failure", DefaultValue: 0}
MqtMsgContentType = Flag{Type: String, Name: flagkey.MqtMsgContentType, Short: "c", Usage: "Content type of messages that publish to the topic", DefaultValue: "application/json"}
MqtPollingInterval = Flag{Type: Int, Name: flagkey.MqtPollingInterval, Usage: "Interval to check the message source for up/down scaling operation of consumers", DefaultValue: 30}
MqtCooldownPeriod = Flag{Type: Int, Name: flagkey.MqtCooldownPeriod, Usage: "The period to wait after the last trigger reported active before scaling the consumer back to 0", DefaultValue: 300}
MqtMinReplicaCount = Flag{Type: Int, Name: flagkey.MqtMinReplicaCount, Usage: "Minimum number of replicas of consumers to scale down to", DefaultValue: 0}
MqtMaxReplicaCount = Flag{Type: Int, Name: flagkey.MqtMaxReplicaCount, Usage: "Maximum number of replicas of consumers to scale up to", DefaultValue: 100}
MqtMetadata = Flag{Type: StringSlice, Name: flagkey.MqtMetadata, Usage: "Metadata needed for connecting to source system in format: --metadata key1=value1 --metadata key2=value2"}
MqtSecret = Flag{Type: String, Name: flagkey.MqtSecret, Usage: "Name of secret object", DefaultValue: ""}
MqtKind = Flag{Type: String, Name: flagkey.MqtKind, Usage: "Kind of Message Queue Trigger, e.g. fission, keda", DefaultValue: "fission"}
EnvName = Flag{Type: String, Name: flagkey.EnvName, Usage: "Environment name"}
EnvPoolsize = Flag{Type: Int, Name: flagkey.EnvPoolsize, Usage: "Size of the pool", DefaultValue: 3}
+15 -8
View File
@@ -81,14 +81,21 @@ const (
TtFnName = "function"
TtRound = "round"
MqtName = resourceName
MqtFnName = "function"
MqtMQType = "mqtype"
MqtTopic = "topic"
MqtRespTopic = "resptopic"
MqtErrorTopic = "errortopic"
MqtMaxRetries = "maxretries"
MqtMsgContentType = "contenttype"
MqtName = resourceName
MqtFnName = "function"
MqtMQType = "mqtype"
MqtTopic = "topic"
MqtRespTopic = "resptopic"
MqtErrorTopic = "errortopic"
MqtMaxRetries = "maxretries"
MqtMsgContentType = "contenttype"
MqtPollingInterval = "pollinginterval"
MqtCooldownPeriod = "cooldownperiod"
MqtMinReplicaCount = "minreplicacount"
MqtMaxReplicaCount = "maxreplicacount"
MqtMetadata = "metadata"
MqtSecret = "secret"
MqtKind = "mqtkind"
EnvName = resourceName
EnvPoolsize = "poolsize"
+15
View File
@@ -315,3 +315,18 @@ func GetSpecDir(input cli.Input) string {
}
return specDir
}
// UpdateMapFromStringSlice parses key, val from "key=val" string array and updates passed map
func UpdateMapFromStringSlice(dataMap *map[string]string, params []string) bool {
updated := false
for _, m := range params {
keyValue := strings.SplitN(m, "=", 2)
if len(keyValue) == 2 {
key := keyValue[0]
value := keyValue[1]
(*dataMap)[key] = value
updated = true
}
}
return updated
}
+19
View File
@@ -0,0 +1,19 @@
FROM golang:1.12-alpine as builder
RUN apk add bash ca-certificates git gcc g++ libc-dev
ARG GOPKG=github.com/fission/fission
ENV GO111MODULE=on
WORKDIR /go/src/${GOPKG}
COPY ./ ./
WORKDIR /go/src/${GOPKG}/pkg/mqtrigger/kafka
RUN go build -a -o /go/bin/main
FROM alpine:3.12 as base
RUN apk add --update ca-certificates
COPY --from=builder /go/bin/main /
ENTRYPOINT ["/main"]
+453
View File
@@ -0,0 +1,453 @@
package main
import (
"context"
"crypto/sha256"
"crypto/sha512"
"crypto/tls"
"crypto/x509"
"fmt"
"hash"
"io/ioutil"
"log"
"net/http"
"os"
"os/signal"
"strings"
"sync"
"syscall"
"time"
"github.com/Shopify/sarama"
"github.com/pkg/errors"
"github.com/xdg/scram"
"go.uber.org/zap"
"github.com/fission/fission/pkg/mqtrigger/util"
)
type kafkaMetadata struct {
bootstrapServers []string
consumerGroup string
// auth
authMode kafkaAuthMode
username string
password string
// ssl
cert string
key string
ca string
}
type kafkaAuthMode string
const (
kafkaAuthModeNone kafkaAuthMode = "none"
kafkaAuthModeSaslPlaintext kafkaAuthMode = "sasl_plaintext"
kafkaAuthModeSaslScramSha256 kafkaAuthMode = "sasl_scram_sha256"
kafkaAuthModeSaslScramSha512 kafkaAuthMode = "sasl_scram_sha512"
kafkaAuthModeSaslSSL kafkaAuthMode = "sasl_ssl"
kafkaAuthModeSaslSSLPlain kafkaAuthMode = "sasl_ssl_plain"
)
var SHA256 scram.HashGeneratorFcn = func() hash.Hash { return sha256.New() }
var SHA512 scram.HashGeneratorFcn = func() hash.Hash { return sha512.New() }
type XDGSCRAMClient struct {
*scram.Client
*scram.ClientConversation
scram.HashGeneratorFcn
}
func (x *XDGSCRAMClient) Begin(userName, password, authzID string) (err error) {
x.Client, err = x.HashGeneratorFcn.NewClient(userName, password, authzID)
if err != nil {
return err
}
x.ClientConversation = x.Client.NewConversation()
return nil
}
func (x *XDGSCRAMClient) Step(challenge string) (response string, err error) {
response, err = x.ClientConversation.Step(challenge)
return
}
func (x *XDGSCRAMClient) Done() bool {
return x.ClientConversation.Done()
}
func parseKafkaMetadata(logger *zap.Logger) (kafkaMetadata, error) {
meta := kafkaMetadata{}
// brokerList marked as deprecated, bootstrapServers is the new one to use
if os.Getenv("BROKER_LIST") != "" && os.Getenv("BOOTSTRAP_SERVERS") != "" {
return meta, errors.New("cannot specify both bootstrapServers and brokerList (deprecated)")
}
if os.Getenv("BROKER_LIST") == "" && os.Getenv("BOOTSTRAP_SERVERS") == "" {
return meta, errors.New("no bootstrapServers or brokerList (deprecated) given")
}
if os.Getenv("BOOTSTRAP_SERVERS") != "" {
meta.bootstrapServers = strings.Split(os.Getenv("BOOTSTRAP_SERVERS"), ",")
}
if os.Getenv("BROKER_LIST") != "" {
logger.Info("WARNING: usage of brokerList is deprecated. use bootstrapServers instead.")
meta.bootstrapServers = strings.Split(os.Getenv("BROKER_LIST"), ",")
}
if os.Getenv("CONSUMER_GROUP") == "" {
return meta, errors.New("No consumerGroup given")
}
meta.consumerGroup = os.Getenv("CONSUMER_GROUP")
meta.authMode = kafkaAuthModeNone
mode := kafkaAuthMode(strings.TrimSpace((os.Getenv("AUTH_MODE"))))
if mode == "" {
mode = kafkaAuthModeNone
}
if mode != kafkaAuthModeNone && mode != kafkaAuthModeSaslPlaintext && mode != kafkaAuthModeSaslSSL && mode != kafkaAuthModeSaslSSLPlain && mode != kafkaAuthModeSaslScramSha256 && mode != kafkaAuthModeSaslScramSha512 {
return meta, fmt.Errorf("err auth mode %s given", mode)
}
meta.authMode = mode
if meta.authMode != kafkaAuthModeNone && meta.authMode != kafkaAuthModeSaslSSL {
if os.Getenv("USERNAME") == "" {
return meta, errors.New("no username given")
}
meta.username = strings.TrimSpace(os.Getenv("USERNAME"))
if os.Getenv("PASSWORD") == "" {
return meta, errors.New("no password given")
}
meta.password = strings.TrimSpace(os.Getenv("PASSWORD"))
}
if meta.authMode == kafkaAuthModeSaslSSL {
if os.Getenv("CA") == "" {
return meta, errors.New("no ca given")
}
meta.ca = os.Getenv("CA")
if os.Getenv("CERT") == "" {
return meta, errors.New("no cert given")
}
meta.cert = os.Getenv("CERT")
if os.Getenv("KEY") == "" {
return meta, errors.New("no key given")
}
meta.key = os.Getenv("KEY")
}
return meta, nil
}
func getConfig(metadata kafkaMetadata) (*sarama.Config, error) {
config := sarama.NewConfig()
config.Version = sarama.V1_0_0_0
if ok := metadata.authMode == kafkaAuthModeSaslPlaintext || metadata.authMode == kafkaAuthModeSaslSSLPlain || metadata.authMode == kafkaAuthModeSaslScramSha256 || metadata.authMode == kafkaAuthModeSaslScramSha512; ok {
config.Net.SASL.Enable = true
config.Net.SASL.User = metadata.username
config.Net.SASL.Password = metadata.password
}
if metadata.authMode == kafkaAuthModeSaslSSLPlain {
config.Net.SASL.Mechanism = sarama.SASLMechanism(sarama.SASLTypePlaintext)
tlsConfig := &tls.Config{
InsecureSkipVerify: true,
ClientAuth: 0,
}
config.Net.TLS.Enable = true
config.Net.TLS.Config = tlsConfig
config.Net.DialTimeout = 10 * time.Second
}
if metadata.authMode == kafkaAuthModeSaslSSL {
cert, err := tls.X509KeyPair([]byte(metadata.cert), []byte(metadata.key))
if err != nil {
return nil, fmt.Errorf("error parse X509KeyPair: %s", err)
}
caCertPool := x509.NewCertPool()
caCertPool.AppendCertsFromPEM([]byte(metadata.ca))
tlsConfig := &tls.Config{
Certificates: []tls.Certificate{cert},
RootCAs: caCertPool,
}
config.Net.TLS.Enable = true
config.Net.TLS.Config = tlsConfig
}
if metadata.authMode == kafkaAuthModeSaslScramSha256 {
config.Net.SASL.SCRAMClientGeneratorFunc = func() sarama.SCRAMClient { return &XDGSCRAMClient{HashGeneratorFcn: SHA256} }
config.Net.SASL.Mechanism = sarama.SASLMechanism(sarama.SASLTypeSCRAMSHA256)
}
if metadata.authMode == kafkaAuthModeSaslScramSha512 {
config.Net.SASL.SCRAMClientGeneratorFunc = func() sarama.SCRAMClient { return &XDGSCRAMClient{HashGeneratorFcn: SHA512} }
config.Net.SASL.Mechanism = sarama.SASLMechanism(sarama.SASLTypeSCRAMSHA512)
}
if metadata.authMode == kafkaAuthModeSaslPlaintext {
config.Net.SASL.Mechanism = sarama.SASLTypePlaintext
config.Net.TLS.Enable = true
}
return config, nil
}
// Connector represents a Sarama consumer group consumer
type Connector struct {
ready chan bool
logger *zap.Logger
producer sarama.SyncProducer
fissionTriggerFields util.FissionMetadata
}
// Setup is run at the beginning of a new session, before ConsumeClaim
func (connector *Connector) Setup(sarama.ConsumerGroupSession) error {
close(connector.ready)
return nil
}
// Cleanup is run at the end of a session, once all ConsumeClaim goroutines have exited
func (connector *Connector) Cleanup(sarama.ConsumerGroupSession) error {
return nil
}
// ConsumeClaim must start a consumer loop of ConsumerGroupClaim's Messages()
func (connector *Connector) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
for message := range claim.Messages() {
connector.logger.Info(fmt.Sprintf("Message claimed: value = %s, timestamp = %v, topic = %s", string(message.Value), message.Timestamp, message.Topic))
success := handleFissionFunction(message, connector.fissionTriggerFields, connector.producer, connector.logger)
if success {
session.MarkMessage(message, "")
}
}
return nil
}
func getProducer(metadata kafkaMetadata) (sarama.SyncProducer, error) {
config, err := getConfig(metadata)
if err != nil {
return nil, err
}
config.Producer.RequiredAcks = sarama.WaitForAll
config.Producer.Retry.Max = 10
config.Producer.Return.Successes = true
producer, err := sarama.NewSyncProducer(metadata.bootstrapServers, config)
if err != nil {
return nil, err
}
return producer, nil
}
func handleFissionFunction(msg *sarama.ConsumerMessage, triggerFields util.FissionMetadata, producer sarama.SyncProducer, logger *zap.Logger) bool {
var value string = string(msg.Value[:])
// Generate the Headers
fissionHeaders := map[string]string{
"X-Fission-MQTrigger-Topic": triggerFields.Topic,
"X-Fission-MQTrigger-RespTopic": triggerFields.ResponseTopic,
"X-Fission-MQTrigger-ErrorTopic": triggerFields.ErrorTopic,
"Content-Type": triggerFields.ContentType,
}
// Create request
req, err := http.NewRequest("POST", triggerFields.FunctionURL, strings.NewReader(value))
if err != nil {
logger.Error("failed to create HTTP request to invoke function",
zap.Error(err),
zap.String("function_url", triggerFields.FunctionURL))
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 <= triggerFields.MaxRetries; attempt++ {
// Make the request
resp, err = http.DefaultClient.Do(req)
if err != nil {
logger.Error("sending function invocation request failed",
zap.Error(err),
zap.String("function_url", triggerFields.FunctionURL),
zap.String("trigger", triggerFields.TriggerName))
continue
}
if resp == nil {
continue
}
if err == nil && resp.StatusCode == http.StatusOK {
// Success, quit retrying
break
}
}
if resp == nil {
logger.Warn("every function invocation retry failed; final retry gave empty response",
zap.String("function_url", triggerFields.FunctionURL),
zap.String("trigger", triggerFields.TriggerName))
return false
}
defer resp.Body.Close()
body, err := ioutil.ReadAll(resp.Body)
logger.Debug("got response from function invocation",
zap.String("function_url", triggerFields.FunctionURL),
zap.String("trigger", triggerFields.TriggerName),
zap.String("body", string(body)))
if err != nil {
errorHandler(logger, triggerFields, producer,
errors.Wrapf(err, "request body error: %v", string(body)))
return false
}
if resp.StatusCode != 200 {
errorHandler(logger, triggerFields, producer,
fmt.Errorf("request returned failure: %v", resp.StatusCode))
return false
}
if len(triggerFields.ResponseTopic) > 0 {
// Generate Kafka record headers
var kafkaRecordHeaders []sarama.RecordHeader
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)})
}
}
_, _, err = producer.SendMessage(&sarama.ProducerMessage{
Topic: triggerFields.ResponseTopic,
Value: sarama.StringEncoder(body),
Headers: kafkaRecordHeaders,
})
if err != nil {
logger.Warn("failed to publish response body from function invocation to topic",
zap.Error(err),
zap.String("topic", triggerFields.Topic),
zap.String("function_url", triggerFields.FunctionURL))
return false
}
}
return true
}
func errorHandler(logger *zap.Logger, triggerFields util.FissionMetadata, producer sarama.SyncProducer, err error) {
if len(triggerFields.ErrorTopic) > 0 {
_, _, e := producer.SendMessage(&sarama.ProducerMessage{
Topic: triggerFields.ErrorTopic,
Value: sarama.StringEncoder(err.Error()),
})
if e != nil {
logger.Error("failed to publish message to error topic",
zap.Error(e),
zap.String("trigger", triggerFields.TriggerName),
zap.String("message", err.Error()),
zap.String("topic", triggerFields.Topic))
}
} else {
logger.Error("message received to publish to error topic, but no error topic was set",
zap.String("message", err.Error()), zap.String("trigger", triggerFields.TriggerName), zap.String("function_url", triggerFields.FunctionURL))
}
}
func main() {
logger, err := zap.NewProduction()
if err != nil {
log.Fatalf("can't initialize zap logger: %v", err)
}
defer logger.Sync()
metadata, err := parseKafkaMetadata(logger)
if err != nil {
logger.Error("Failed to fetch kafka metadata", zap.Error(err))
return
}
triggerFields, err := util.ParseFissionMetadata()
if err != nil {
logger.Error("Failed to parse fission trigger fields", zap.Error(err))
return
}
config, err := getConfig(metadata)
if err != nil {
logger.Error("Failed to create kafka config", zap.Error(err))
return
}
producer, err := getProducer(metadata)
if err != nil {
logger.Error("Failed to create kafka producer", zap.Error(err))
return
}
defer producer.Close()
connector := Connector{
ready: make(chan bool),
logger: logger,
producer: producer,
fissionTriggerFields: triggerFields,
}
ctx, cancel := context.WithCancel(context.Background())
client, err := sarama.NewConsumerGroup(metadata.bootstrapServers, metadata.consumerGroup, config)
if err != nil {
logger.Error("Error creating consumer group client", zap.Error(err))
return
}
wg := &sync.WaitGroup{}
wg.Add(1)
go func() {
defer wg.Done()
for {
if err := client.Consume(ctx, []string{triggerFields.Topic}, &connector); err != nil {
logger.Error("Error from consumer", zap.Error(err))
}
// check if context was cancelled, signaling that the consumer should stop
if ctx.Err() != nil {
return
}
connector.ready = make(chan bool)
}
}()
<-connector.ready // Await till the consumer has been set up
logger.Info("Sarama consumer up and running!...")
sigterm := make(chan os.Signal, 1)
signal.Notify(sigterm, syscall.SIGINT, syscall.SIGTERM)
select {
case <-ctx.Done():
logger.Info("terminating: context cancelled")
case <-sigterm:
logger.Info("terminating: via signal")
}
cancel()
wg.Wait()
if err = client.Close(); err != nil {
logger.Error("Error closing client", zap.Error(err))
}
}
+540
View File
@@ -0,0 +1,540 @@
package mqtrigger
import (
"context"
"fmt"
"os"
"regexp"
"strconv"
"strings"
"time"
"github.com/pkg/errors"
"go.uber.org/zap"
appsv1 "k8s.io/api/apps/v1"
apiv1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/fields"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/kubernetes"
k8sCache "k8s.io/client-go/tools/cache"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/fission/fission/pkg/crd"
"github.com/fission/fission/pkg/utils"
)
var (
scaledObjectGVR = schema.GroupVersionResource{
Group: "keda.k8s.io",
Version: "v1alpha1",
Resource: "scaledobjects",
}
authTriggerGVR = schema.GroupVersionResource{
Group: "keda.k8s.io",
Version: "v1alpha1",
Resource: "triggerauthentications",
}
matchFirstCap = regexp.MustCompile("(.)([A-Z][a-z]+)")
matchAllCap = regexp.MustCompile("([a-z0-9])([A-Z])")
)
func getScaledObjectClient(namespace string) (dynamic.ResourceInterface, error) {
dynamicClient, err := crd.GetDynamicClient()
if err != nil {
return nil, err
}
return dynamicClient.Resource(scaledObjectGVR).Namespace(namespace), nil
}
func getAuthTriggerClient(namespace string) (dynamic.ResourceInterface, error) {
dynamicClient, err := crd.GetDynamicClient()
if err != nil {
return nil, err
}
return dynamicClient.Resource(scaledObjectGVR).Namespace(namespace), nil
}
// StartScalerManager watches for changes in MessageQueueTrigger and,
// Based on changes, it Creates, Updates and Deletes Objects of Kind ScaledObjects, AuthenticationTriggers and Deployments
func StartScalerManager(logger *zap.Logger, routerURL string) error {
fissionClient, kubeClient, _, err := crd.MakeFissionClient()
if err != nil {
return err
}
err = fissionClient.WaitForCRDs()
if err != nil {
return errors.Wrap(err, "error waiting for CRDs")
}
crdClient := fissionClient.CoreV1().RESTClient()
resyncPeriod := 30 * time.Second
listWatch := k8sCache.NewListWatchFromClient(crdClient, "messagequeuetriggers", metav1.NamespaceAll, fields.Everything())
_, controller := k8sCache.NewInformer(listWatch, &fv1.MessageQueueTrigger{}, resyncPeriod, k8sCache.ResourceEventHandlerFuncs{
AddFunc: func(obj interface{}) {
go func() {
mqt := obj.(*fv1.MessageQueueTrigger)
if mqt.Spec.MqtKind == "fission" {
return
}
logger.Debug("Create deployment for Scaler Object", zap.Any("mqt", mqt.ObjectMeta), zap.Any("mqt.Spec", mqt.Spec))
authenticationRef := ""
if len(mqt.Spec.Secret) > 0 {
authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name)
err = createAuthTrigger(mqt, authenticationRef, kubeClient)
if err != nil {
logger.Error("Failed to create Authentication Trigger", zap.Error(err))
return
}
}
if err = createDeployment(mqt, routerURL, kubeClient); err != nil {
logger.Error("Failed to create Deployment", zap.Error(err))
if len(authenticationRef) > 0 {
err = deleteAuthTrigger(authenticationRef, mqt.ObjectMeta.Namespace)
if err != nil {
logger.Error("Failed to delete Authentication Trigger", zap.Error(err))
}
}
return
}
if err = createScaledObject(mqt, authenticationRef); err != nil {
logger.Error("Failed to create ScaledObject", zap.Error(err))
if len(authenticationRef) > 0 {
if err = deleteAuthTrigger(authenticationRef, mqt.ObjectMeta.Namespace); err != nil {
logger.Error("Failed to delete Authentication Trigger", zap.Error(err))
}
}
if err = deleteDeployment(mqt.ObjectMeta.Name, kubeClient); err != nil {
logger.Error("Failed to delete Deployment", zap.Error(err))
}
}
}()
},
UpdateFunc: func(obj interface{}, newObj interface{}) {
go func() {
mqt := obj.(*fv1.MessageQueueTrigger)
newMqt := newObj.(*fv1.MessageQueueTrigger)
updated := checkAndUpdateTriggerFields(mqt, newMqt)
if mqt.Spec.MqtKind == "fission" {
return
}
if !updated {
logger.Warn(fmt.Sprintf("%s remains unchanged. No changes found in trigger fields", mqt.ObjectMeta.Name))
return
}
authenticationRef := ""
if len(newMqt.Spec.Secret) > 0 && newMqt.Spec.Secret != mqt.Spec.Secret {
authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name)
if err = updateAuthTrigger(mqt, authenticationRef, kubeClient); err != nil {
logger.Error("Failed to update Authentication Trigger", zap.Error(err))
return
}
}
if err = updateDeployment(mqt, routerURL, kubeClient); err != nil {
logger.Error("Failed to Update Deployment", zap.Error(err))
return
}
if err = updateScaledObject(mqt, authenticationRef); err != nil {
logger.Error("Failed to Update ScaledObject", zap.Error(err))
return
}
}()
},
})
controller.Run(context.Background().Done())
return nil
}
func toEnvVar(str string) string {
envVar := matchFirstCap.ReplaceAllString(str, "${1}_${2}")
envVar = matchAllCap.ReplaceAllString(envVar, "${1}_${2}")
return strings.ToUpper(envVar)
}
func getEnvVarlist(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient kubernetes.Interface) ([]apiv1.EnvVar, error) {
url := routerURL + "/" + strings.TrimPrefix(utils.UrlForFunction(mqt.Spec.FunctionReference.Name, mqt.ObjectMeta.Namespace), "/")
envVars := []apiv1.EnvVar{
{
Name: "TOPIC",
Value: mqt.Spec.Topic,
},
{
Name: "FUNCTION_URL",
Value: url,
},
{
Name: "ERROR_TOPIC",
Value: mqt.Spec.ErrorTopic,
},
{
Name: "RESPONSE_TOPIC",
Value: mqt.Spec.ResponseTopic,
},
{
Name: "TRIGGER_NAME",
Value: mqt.ObjectMeta.Name,
},
{
Name: "MAX_RETRIES",
Value: strconv.Itoa(mqt.Spec.MaxRetries),
},
{
Name: "CONTENT_TYPE",
Value: mqt.Spec.ContentType,
},
}
// Metadata Fields
for key, value := range mqt.Spec.Metadata {
envVars = append(envVars, apiv1.EnvVar{
Name: toEnvVar(key),
Value: value,
})
}
// Add Auth Fields
secretName := mqt.Spec.Secret
if len(secretName) > 0 {
secret, err := kubeClient.CoreV1().Secrets(apiv1.NamespaceDefault).Get(secretName, metav1.GetOptions{})
if err != nil {
return nil, err
}
for key, value := range secret.Data {
envVars = append(envVars, apiv1.EnvVar{
Name: toEnvVar(key),
Value: string(value),
})
}
}
return envVars, nil
}
func checkAndUpdateTriggerFields(mqt, newMqt *fv1.MessageQueueTrigger) bool {
updated := false
if len(newMqt.Spec.Topic) > 0 && newMqt.Spec.Topic != mqt.Spec.Topic {
mqt.Spec.Topic = newMqt.Spec.Topic
updated = true
}
if len(newMqt.Spec.ResponseTopic) > 0 && newMqt.Spec.ResponseTopic != mqt.Spec.ResponseTopic {
mqt.Spec.ResponseTopic = newMqt.Spec.ResponseTopic
updated = true
}
if len(newMqt.Spec.ErrorTopic) > 0 && newMqt.Spec.ErrorTopic != mqt.Spec.ErrorTopic {
mqt.Spec.ErrorTopic = newMqt.Spec.ErrorTopic
updated = true
}
if newMqt.Spec.MaxRetries >= 0 && newMqt.Spec.MaxRetries != mqt.Spec.MaxRetries {
mqt.Spec.MaxRetries = newMqt.Spec.MaxRetries
updated = true
}
if len(newMqt.Spec.FunctionReference.Name) > 0 && newMqt.Spec.FunctionReference.Name != mqt.Spec.FunctionReference.Name {
mqt.Spec.FunctionReference.Name = newMqt.Spec.FunctionReference.Name
updated = true
}
if len(newMqt.Spec.ContentType) > 0 && newMqt.Spec.ContentType != mqt.Spec.ContentType {
mqt.Spec.ContentType = newMqt.Spec.ContentType
updated = true
}
if *newMqt.Spec.PollingInterval >= 0 && *newMqt.Spec.PollingInterval != *mqt.Spec.PollingInterval {
mqt.Spec.PollingInterval = newMqt.Spec.PollingInterval
updated = true
}
if *newMqt.Spec.CooldownPeriod >= 0 && *newMqt.Spec.CooldownPeriod != *mqt.Spec.CooldownPeriod {
mqt.Spec.CooldownPeriod = newMqt.Spec.CooldownPeriod
updated = true
}
if *newMqt.Spec.MinReplicaCount >= 0 && *newMqt.Spec.MinReplicaCount != *mqt.Spec.MinReplicaCount {
mqt.Spec.MinReplicaCount = newMqt.Spec.MinReplicaCount
updated = true
}
if *newMqt.Spec.MaxReplicaCount >= 0 && *newMqt.Spec.MaxReplicaCount != *mqt.Spec.MaxReplicaCount {
mqt.Spec.MaxReplicaCount = newMqt.Spec.MaxReplicaCount
updated = true
}
if len(newMqt.Spec.FunctionReference.Name) > 0 && newMqt.Spec.FunctionReference.Name != mqt.Spec.FunctionReference.Name {
newMqt.Spec.FunctionReference.Name = mqt.Spec.FunctionReference.Name
updated = true
}
for key, value := range newMqt.Spec.Metadata {
if val, ok := mqt.Spec.Metadata[key]; ok && val != value {
mqt.Spec.Metadata[key] = value
updated = true
}
}
if len(newMqt.Spec.Secret) > 0 && newMqt.Spec.Secret != mqt.Spec.Secret {
mqt.Spec.Secret = newMqt.Spec.Secret
updated = true
}
if newMqt.Spec.MqtKind != mqt.Spec.MqtKind {
mqt.Spec.MqtKind = newMqt.Spec.MqtKind
updated = true
}
return updated
}
func getResourceVersion(scaledObjectName string, kedaClient dynamic.ResourceInterface) (version string, err error) {
scaledObject, err := kedaClient.Get(scaledObjectName, metav1.GetOptions{})
if err != nil {
return "", err
}
return scaledObject.GetResourceVersion(), nil
}
func getAuthTriggerSpec(mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient kubernetes.Interface) (*unstructured.Unstructured, error) {
secret, err := kubeClient.CoreV1().Secrets(apiv1.NamespaceDefault).Get(mqt.Spec.Secret, metav1.GetOptions{})
if err != nil {
return nil, err
}
var secretTargetRefFields []interface{}
for secretField := range secret.Data {
secretTargetRefFields = append(secretTargetRefFields, map[string]interface{}{
"name": mqt.Spec.Secret,
"parameter": secretField,
"key": secretField,
})
}
authTriggerObj := &unstructured.Unstructured{
Object: map[string]interface{}{
"kind": "TriggerAuthentication",
"apiVersion": "keda.k8s.io/v1alpha1",
"metadata": map[string]interface{}{
"name": authenticationRef,
"namespace": mqt.ObjectMeta.Namespace,
"ownerReferences": []interface{}{
map[string]interface{}{
"kind": "MessageQueueTrigger",
"apiVersion": "fission.io/v1",
"name": mqt.ObjectMeta.Name,
"uid": mqt.ObjectMeta.UID,
"blockOwnerDeletion": true,
},
},
},
"spec": map[string]interface{}{
"secretTargetRef": secretTargetRefFields,
},
},
}
return authTriggerObj, nil
}
func createAuthTrigger(mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient *kubernetes.Clientset) error {
authTriggerObj, err := getAuthTriggerSpec(mqt, authenticationRef, kubeClient)
if err != nil {
return err
}
authTriggerClient, err := getAuthTriggerClient(mqt.ObjectMeta.Namespace)
if err != nil {
return err
}
_, err = authTriggerClient.Create(authTriggerObj, metav1.CreateOptions{})
if err != nil {
return err
}
return nil
}
func updateAuthTrigger(mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient *kubernetes.Clientset) error {
authTriggerClient, err := getAuthTriggerClient(mqt.ObjectMeta.Namespace)
if err != nil {
return err
}
oldAuthTriggerObj, err := authTriggerClient.Get(authenticationRef, metav1.GetOptions{})
if err != nil {
return err
}
resourceVersion := oldAuthTriggerObj.GetResourceVersion()
authTriggerObj, err := getAuthTriggerSpec(mqt, authenticationRef, kubeClient)
if err != nil {
return err
}
authTriggerObj.SetResourceVersion(resourceVersion)
_, err = authTriggerClient.Update(authTriggerObj, metav1.UpdateOptions{})
if err != nil {
return err
}
return nil
}
func deleteAuthTrigger(name, namespace string) error {
authTriggerClient, err := getAuthTriggerClient(namespace)
if err != nil {
return err
}
err = authTriggerClient.Delete(name, &metav1.DeleteOptions{})
if err != nil {
return err
}
return nil
}
func getDeploymentSpec(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient *kubernetes.Clientset) (*appsv1.Deployment, error) {
envVars, err := getEnvVarlist(mqt, routerURL, kubeClient)
if err != nil {
return nil, err
}
imageName := fmt.Sprintf("%s_image", string(mqt.Spec.MessageQueueType))
image := os.Getenv(strings.ToUpper(imageName))
blockOwnerDeletion := true
return &appsv1.Deployment{
ObjectMeta: metav1.ObjectMeta{
Name: mqt.ObjectMeta.Name,
Labels: map[string]string{
"app": mqt.ObjectMeta.Name,
},
OwnerReferences: []metav1.OwnerReference{
{
Kind: "MessageQueueTrigger",
APIVersion: "fission.io/v1",
Name: mqt.ObjectMeta.Name,
UID: mqt.ObjectMeta.UID,
BlockOwnerDeletion: &blockOwnerDeletion,
},
},
},
Spec: appsv1.DeploymentSpec{
Selector: &metav1.LabelSelector{
MatchLabels: map[string]string{
"app": mqt.ObjectMeta.Name,
},
},
Template: apiv1.PodTemplateSpec{
ObjectMeta: metav1.ObjectMeta{
Labels: map[string]string{
"app": mqt.ObjectMeta.Name,
},
},
Spec: apiv1.PodSpec{
Containers: []apiv1.Container{
{
Name: mqt.ObjectMeta.Name,
Image: image,
ImagePullPolicy: "Always",
Env: envVars,
},
},
},
},
},
}, nil
}
func createDeployment(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient *kubernetes.Clientset) error {
deployment, err := getDeploymentSpec(mqt, routerURL, kubeClient)
if err != nil {
return err
}
_, err = kubeClient.AppsV1().Deployments(apiv1.NamespaceDefault).Create(deployment)
if err != nil {
return err
}
return nil
}
func updateDeployment(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient *kubernetes.Clientset) error {
deployment, err := getDeploymentSpec(mqt, routerURL, kubeClient)
if err != nil {
return err
}
_, err = kubeClient.AppsV1().Deployments(apiv1.NamespaceDefault).Update(deployment)
if err != nil {
return err
}
return nil
}
func deleteDeployment(name string, kubeClient *kubernetes.Clientset) error {
deletePolicy := metav1.DeletePropagationForeground
if err := kubeClient.AppsV1().Deployments(apiv1.NamespaceDefault).Delete(name, &metav1.DeleteOptions{
PropagationPolicy: &deletePolicy,
}); err != nil {
return err
}
return nil
}
func getScaledObject(mqt *fv1.MessageQueueTrigger, authenticationRef string) *unstructured.Unstructured {
return &unstructured.Unstructured{
Object: map[string]interface{}{
"kind": "ScaledObject",
"apiVersion": "keda.k8s.io/v1alpha1",
"metadata": map[string]interface{}{
"name": mqt.ObjectMeta.Name,
"namespace": mqt.ObjectMeta.Namespace,
"ownerReferences": []interface{}{
map[string]interface{}{
"kind": "MessageQueueTrigger",
"apiVersion": "fission.io/v1",
"name": mqt.ObjectMeta.Name,
"uid": mqt.ObjectMeta.UID,
"blockOwnerDeletion": true,
},
},
},
"spec": map[string]interface{}{
"cooldownPeriod": &mqt.Spec.CooldownPeriod,
"maxReplicaCount": &mqt.Spec.MaxReplicaCount,
"minReplicaCount": &mqt.Spec.MinReplicaCount,
"pollingInterval": &mqt.Spec.PollingInterval,
"scaleTargetRef": map[string]interface{}{
"deploymentName": mqt.ObjectMeta.Name,
},
"triggers": []interface{}{
map[string]interface{}{
"type": mqt.Spec.MessageQueueType,
"metadata": mqt.Spec.Metadata,
"authenticationRef": map[string]interface{}{
"name": authenticationRef,
},
},
},
},
},
}
}
func createScaledObject(mqt *fv1.MessageQueueTrigger, authenticationRef string) error {
scaledObject := getScaledObject(mqt, authenticationRef)
kedaClient, err := getScaledObjectClient(mqt.ObjectMeta.Namespace)
if err != nil {
return err
}
_, err = kedaClient.Create(scaledObject, metav1.CreateOptions{})
if err != nil {
return err
}
return nil
}
func updateScaledObject(mqt *fv1.MessageQueueTrigger, authenticationRef string) error {
kedaClient, err := getScaledObjectClient(mqt.ObjectMeta.Namespace)
if err != nil {
return err
}
oldScaledObject, err := kedaClient.Get(mqt.ObjectMeta.Name, metav1.GetOptions{})
if err != nil {
return err
}
resourceVersion := oldScaledObject.GetResourceVersion()
scaledObject := getScaledObject(mqt, authenticationRef)
scaledObject.SetResourceVersion(resourceVersion)
_, err = kedaClient.Update(scaledObject, metav1.UpdateOptions{})
if err != nil {
return err
}
return nil
}
+622
View File
@@ -0,0 +1,622 @@
package mqtrigger
import (
"fmt"
"reflect"
"sort"
"testing"
fv1 "github.com/fission/fission/pkg/apis/core/v1"
"github.com/stretchr/testify/assert"
apiv1 "k8s.io/api/core/v1"
v1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
"k8s.io/apimachinery/pkg/runtime"
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/client-go/dynamic"
dynfake "k8s.io/client-go/dynamic/fake"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/kubernetes/fake"
)
func Test_toEnvVar(t *testing.T) {
type args struct {
str string
}
tests := []struct {
name string
args args
want string
}{
{"Empty string", args{""}, ""},
{"Single word", args{"fission"}, "FISSION"},
{"CamelCase", args{"responseTopic"}, "RESPONSE_TOPIC"},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := toEnvVar(tt.args.str); got != tt.want {
t.Errorf("toEnvVar() = %v, want %v", got, tt.want)
}
})
}
}
func Test_getEnvVarlist(t *testing.T) {
// Kafka Test with Valid Secret
pollingInterval := int32(30)
cooldownPeriod := int32(300)
minReplicaCount := int32(0)
maxReplicaCount := int32(100)
mqt := &fv1.MessageQueueTrigger{
ObjectMeta: metav1.ObjectMeta{
Name: "Test",
Namespace: "default",
},
Spec: fv1.MessageQueueTriggerSpec{
FunctionReference: fv1.FunctionReference{
Type: fv1.FunctionReferenceTypeFunctionName,
Name: "test",
},
MessageQueueType: "kafka",
Topic: "topic",
ResponseTopic: "response-topic",
ErrorTopic: "error-topic",
MaxRetries: 4,
ContentType: "application/json",
PollingInterval: &pollingInterval,
CooldownPeriod: &cooldownPeriod,
MinReplicaCount: &minReplicaCount,
MaxReplicaCount: &maxReplicaCount,
Metadata: map[string]string{
"bootstrapServers": "my-cluster-kafka-brokers.my-kafka-project.svc:9092",
"consumerGroup": "my-group",
"topic": "topic",
},
Secret: "test-kafka-secrets",
MqtKind: "keda",
},
}
data := map[string][]byte{
"authMode": []byte("sasl_plaintext"),
"username": []byte("admin"),
"password": []byte("admin"),
"ca": []byte("test_ca"),
"cert": []byte("test_cert"),
"key": []byte("test_key"),
}
namespace := apiv1.NamespaceDefault
routerURL := "http://router.fission/fission-function"
secret := &v1.Secret{
ObjectMeta: metav1.ObjectMeta{
Name: "test-kafka-secrets",
Namespace: namespace,
},
Data: data,
}
kubeClient := fake.NewSimpleClientset()
_, err := kubeClient.CoreV1().Secrets(namespace).Create(secret)
if err != nil {
assert.Equal(t, nil, err)
}
expectedEnvVars := []apiv1.EnvVar{
{
Name: "TOPIC",
Value: mqt.Spec.Topic,
},
{
Name: "FUNCTION_URL",
Value: "http://router.fission/fission-function/fission-function/test",
},
{
Name: "ERROR_TOPIC",
Value: "error-topic",
},
{
Name: "RESPONSE_TOPIC",
Value: "response-topic",
},
{
Name: "TRIGGER_NAME",
Value: "Test",
},
{
Name: "MAX_RETRIES",
Value: "4",
},
{
Name: "CONTENT_TYPE",
Value: "application/json",
},
{
Name: "BOOTSTRAP_SERVERS",
Value: "my-cluster-kafka-brokers.my-kafka-project.svc:9092",
},
{
Name: "CONSUMER_GROUP",
Value: "my-group",
},
{
Name: "TOPIC",
Value: "topic",
},
{
Name: "KEY",
Value: "test_key",
},
{
Name: "AUTH_MODE",
Value: "sasl_plaintext",
},
{
Name: "USERNAME",
Value: "admin",
},
{
Name: "PASSWORD",
Value: "admin",
},
{
Name: "CA",
Value: "test_ca",
},
{
Name: "CERT",
Value: "test_cert",
},
}
// Kafka Test with Invalid Secret Name
mqt2 := &fv1.MessageQueueTrigger{
ObjectMeta: metav1.ObjectMeta{
Name: "Test2",
Namespace: "default",
},
Spec: fv1.MessageQueueTriggerSpec{
FunctionReference: fv1.FunctionReference{
Type: fv1.FunctionReferenceTypeFunctionName,
Name: "test2",
},
MessageQueueType: "kafka",
Topic: "topic",
ResponseTopic: "response-topic",
ErrorTopic: "error-topic",
MaxRetries: 4,
ContentType: "application/json",
PollingInterval: &pollingInterval,
CooldownPeriod: &cooldownPeriod,
MinReplicaCount: &minReplicaCount,
MaxReplicaCount: &maxReplicaCount,
Metadata: map[string]string{
"bootstrapServers": "my-cluster-kafka-brokers.my-kafka-project.svc:9092",
"consumerGroup": "my-group",
"topic": "topic",
},
Secret: "test-kafka-secrets-invalid",
MqtKind: "keda",
},
}
// Test Code
type args struct {
mqt *fv1.MessageQueueTrigger
routerURL string
kubeClient kubernetes.Interface
}
tests := []struct {
name string
args args
want []apiv1.EnvVar
wantErr bool
}{
{"Test kafka example", args{mqt, routerURL, kubeClient}, expectedEnvVars, false},
{"Test kafka invalid secret", args{mqt2, routerURL, kubeClient}, nil, true},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, err := getEnvVarlist(tt.args.mqt, tt.args.routerURL, tt.args.kubeClient)
sort.Slice(got, func(i, j int) bool {
return got[i].Name < got[j].Name
})
sort.Slice(tt.want, func(i, j int) bool {
return tt.want[i].Name < tt.want[j].Name
})
if (err != nil) != tt.wantErr {
t.Errorf("getEnvVarlist() error = %v, wantErr %v", err, tt.wantErr)
return
}
if !reflect.DeepEqual(got, tt.want) {
t.Errorf("getEnvVarlist() = %v, want %v", got, tt.want)
}
})
}
}
func Test_checkAndUpdateTriggerFields(t *testing.T) {
pollingInterval := int32(30)
cooldownPeriod := int32(300)
minReplicaCount := int32(0)
maxReplicaCount := int32(100)
// Test 1 with difference
mqt := &fv1.MessageQueueTrigger{
ObjectMeta: metav1.ObjectMeta{
Name: "Test",
Namespace: "default",
},
Spec: fv1.MessageQueueTriggerSpec{
FunctionReference: fv1.FunctionReference{
Type: fv1.FunctionReferenceTypeFunctionName,
Name: "test",
},
MessageQueueType: "kafka",
Topic: "topic",
ResponseTopic: "response-topic",
ErrorTopic: "error-topic",
MaxRetries: 4,
ContentType: "application/json",
PollingInterval: &pollingInterval,
CooldownPeriod: &cooldownPeriod,
MinReplicaCount: &minReplicaCount,
MaxReplicaCount: &maxReplicaCount,
Metadata: map[string]string{
"bootstrapServers": "my-cluster-kafka-brokers.my-kafka-project.svc:9092",
"consumerGroup": "my-group",
"topic": "topic",
},
Secret: "test-kafka-secrets",
MqtKind: "keda",
},
}
newMqt1 := &fv1.MessageQueueTrigger{
ObjectMeta: metav1.ObjectMeta{
Name: "Test",
Namespace: "default",
},
Spec: fv1.MessageQueueTriggerSpec{
FunctionReference: fv1.FunctionReference{
Type: fv1.FunctionReferenceTypeFunctionName,
Name: "test2",
},
MessageQueueType: "kafka",
Topic: "my-topic",
ResponseTopic: "response-topic",
ErrorTopic: "error-topic",
MaxRetries: 2,
ContentType: "application/json",
PollingInterval: &pollingInterval,
CooldownPeriod: &cooldownPeriod,
MinReplicaCount: &minReplicaCount,
MaxReplicaCount: &maxReplicaCount,
Metadata: map[string]string{
"bootstrapServers": "my-cluster-kafka-brokers-2.my-kafka-project.svc:9092",
"consumerGroup": "my-group-2",
"topic": "my-topic",
},
Secret: "new-test-kafka-secrets",
MqtKind: "keda",
},
}
// Test 2 with no difference
mqt2 := &fv1.MessageQueueTrigger{
ObjectMeta: metav1.ObjectMeta{
Name: "Test",
Namespace: "default",
},
Spec: fv1.MessageQueueTriggerSpec{
FunctionReference: fv1.FunctionReference{
Type: fv1.FunctionReferenceTypeFunctionName,
Name: "test",
},
MessageQueueType: "kafka",
Topic: "topic",
ResponseTopic: "response-topic",
ErrorTopic: "error-topic",
MaxRetries: 4,
ContentType: "application/json",
PollingInterval: &pollingInterval,
CooldownPeriod: &cooldownPeriod,
MinReplicaCount: &minReplicaCount,
MaxReplicaCount: &maxReplicaCount,
Metadata: map[string]string{
"bootstrapServers": "my-cluster-kafka-brokers.my-kafka-project.svc:9092",
"consumerGroup": "my-group",
"topic": "topic",
},
Secret: "test-kafka-secrets",
MqtKind: "keda",
},
}
newMqt2 := &fv1.MessageQueueTrigger{
ObjectMeta: metav1.ObjectMeta{
Name: "Test",
Namespace: "default",
},
Spec: fv1.MessageQueueTriggerSpec{
FunctionReference: fv1.FunctionReference{
Type: fv1.FunctionReferenceTypeFunctionName,
Name: "test",
},
MaxRetries: 4,
PollingInterval: &pollingInterval,
CooldownPeriod: &cooldownPeriod,
MinReplicaCount: &minReplicaCount,
MaxReplicaCount: &maxReplicaCount,
MqtKind: "keda",
},
}
type args struct {
mqt *fv1.MessageQueueTrigger
newMqt *fv1.MessageQueueTrigger
}
tests := []struct {
name string
args args
want bool
}{
{"With diff", args{mqt, newMqt1}, true},
{"With no diff", args{mqt2, newMqt2}, false},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
if got := checkAndUpdateTriggerFields(tt.args.mqt, tt.args.newMqt); got != tt.want {
t.Errorf("checkAndUpdateTriggerFields() = %v, want %v", got, tt.want)
}
})
}
}
func newUnstructured(apiVersion, kind, namespace, name, resourceVersion string) *unstructured.Unstructured {
return &unstructured.Unstructured{
Object: map[string]interface{}{
"apiVersion": apiVersion,
"kind": kind,
"metadata": map[string]interface{}{
"namespace": namespace,
"name": name,
"resourceVersion": resourceVersion,
},
},
}
}
func Test_getResourceVersion(t *testing.T) {
scheme := runtime.NewScheme()
client := dynfake.NewSimpleDynamicClient(scheme, newUnstructured("keda.k8s.io/v1alpha1", "ScaledObject", "default", "test-1", "12345"))
dynamicResourceClient := client.Resource(schema.GroupVersionResource{
Group: "keda.k8s.io",
Version: "v1alpha1",
Resource: "scaledobjects",
})
type args struct {
scaledObjectName string
kedaClient dynamic.ResourceInterface
}
tests := []struct {
name string
args args
wantVersion string
wantErr bool
}{
{"Valid Resource", args{"test-1", dynamicResourceClient}, "12345", false},
{"Invalid Resource", args{"test-2", dynamicResourceClient}, "", true},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
gotVersion, err := getResourceVersion(tt.args.scaledObjectName, tt.args.kedaClient)
if (err != nil) != tt.wantErr {
t.Errorf("getResourceVersion() error = %v, wantErr %v", err, tt.wantErr)
return
}
if gotVersion != tt.wantVersion {
t.Errorf("getResourceVersion() = %v, want %v", gotVersion, tt.wantVersion)
}
})
}
}
func Test_getAuthTriggerSpec(t *testing.T) {
// Valid - with Secret
pollingInterval := int32(30)
cooldownPeriod := int32(300)
minReplicaCount := int32(0)
maxReplicaCount := int32(200)
mqt1 := &fv1.MessageQueueTrigger{
ObjectMeta: metav1.ObjectMeta{
Name: "Test",
Namespace: "default",
UID: "test123",
},
Spec: fv1.MessageQueueTriggerSpec{
FunctionReference: fv1.FunctionReference{
Type: fv1.FunctionReferenceTypeFunctionName,
Name: "test",
},
MessageQueueType: "kafka",
Topic: "topic",
ResponseTopic: "response-topic",
ErrorTopic: "error-topic",
MaxRetries: 4,
ContentType: "application/json",
PollingInterval: &pollingInterval,
CooldownPeriod: &cooldownPeriod,
MinReplicaCount: &minReplicaCount,
MaxReplicaCount: &maxReplicaCount,
Metadata: map[string]string{
"bootstrapServers": "my-cluster-kafka-brokers.my-kafka-project.svc:9092",
"consumerGroup": "my-group",
"topic": "topic",
},
Secret: "test-kafka-secrets",
MqtKind: "keda",
},
}
data := map[string][]byte{
"authMode": []byte("sasl_plaintext"),
"username": []byte("admin"),
"password": []byte("admin"),
"ca": []byte("test_ca"),
"cert": []byte("test_cert"),
"key": []byte("test_key"),
}
namespace := apiv1.NamespaceDefault
secret := &v1.Secret{
ObjectMeta: metav1.ObjectMeta{
Name: "test-kafka-secrets",
Namespace: namespace,
},
Data: data,
}
kubeClient := fake.NewSimpleClientset()
_, err := kubeClient.CoreV1().Secrets(namespace).Create(secret)
if err != nil {
assert.Equal(t, nil, err)
}
authenticationRef := fmt.Sprintf("%s-auth-trigger", mqt1.ObjectMeta.Name)
expectedAuthTriggerObj := &unstructured.Unstructured{
Object: map[string]interface{}{
"kind": "TriggerAuthentication",
"apiVersion": "keda.k8s.io/v1alpha1",
"metadata": map[string]interface{}{
"name": authenticationRef,
"namespace": mqt1.ObjectMeta.Namespace,
"ownerReferences": []interface{}{
map[string]interface{}{
"kind": "MessageQueueTrigger",
"apiVersion": "fission.io/v1",
"name": mqt1.ObjectMeta.Name,
"uid": mqt1.ObjectMeta.UID,
"blockOwnerDeletion": true,
},
},
},
"spec": map[string]interface{}{
"secretTargetRef": []interface{}{
map[string]interface{}{
"name": mqt1.Spec.Secret,
"parameter": "authMode",
"key": "authMode",
},
map[string]interface{}{
"name": mqt1.Spec.Secret,
"parameter": "username",
"key": "username",
},
map[string]interface{}{
"name": mqt1.Spec.Secret,
"parameter": "password",
"key": "password",
},
map[string]interface{}{
"name": mqt1.Spec.Secret,
"parameter": "ca",
"key": "ca",
},
map[string]interface{}{
"name": mqt1.Spec.Secret,
"parameter": "cert",
"key": "cert",
},
map[string]interface{}{
"name": mqt1.Spec.Secret,
"parameter": "key",
"key": "key",
},
},
},
},
}
// Invalid without secret
mqt2 := &fv1.MessageQueueTrigger{
ObjectMeta: metav1.ObjectMeta{
Name: "Test",
Namespace: "default",
UID: "test123",
},
Spec: fv1.MessageQueueTriggerSpec{
FunctionReference: fv1.FunctionReference{
Type: fv1.FunctionReferenceTypeFunctionName,
Name: "test",
},
MessageQueueType: "kafka",
Topic: "topic",
ResponseTopic: "response-topic",
ErrorTopic: "error-topic",
MaxRetries: 4,
ContentType: "application/json",
PollingInterval: &pollingInterval,
CooldownPeriod: &cooldownPeriod,
MinReplicaCount: &minReplicaCount,
MaxReplicaCount: &maxReplicaCount,
Metadata: map[string]string{
"bootstrapServers": "my-cluster-kafka-brokers.my-kafka-project.svc:9092",
"consumerGroup": "my-group",
"topic": "topic",
},
Secret: "test-kafka-no-secret",
MqtKind: "keda",
},
}
type args struct {
mqt *fv1.MessageQueueTrigger
authenticationRef string
kubeClient kubernetes.Interface
}
tests := []struct {
name string
args args
want *unstructured.Unstructured
wantErr bool
}{
{"With secret", args{mqt1, authenticationRef, kubeClient}, expectedAuthTriggerObj, false},
{"With invalid secret", args{mqt2, authenticationRef, kubeClient}, nil, true},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
got, err := getAuthTriggerSpec(tt.args.mqt, tt.args.authenticationRef, tt.args.kubeClient)
if (err != nil) != tt.wantErr {
t.Errorf("getAuthTriggerSpec() error = %v, wantErr %v", err, tt.wantErr)
return
}
if err != nil && tt.wantErr {
return
}
gotSpec := got.Object["spec"].(map[string]interface{})["secretTargetRef"].([]interface{})
sort.Slice(gotSpec, func(i, j int) bool {
return gotSpec[i].(map[string]interface{})["parameter"].(string) < gotSpec[j].(map[string]interface{})["parameter"].(string)
})
wantSpec := tt.want.Object["spec"].(map[string]interface{})["secretTargetRef"].([]interface{})
sort.Slice(wantSpec, func(i, j int) bool {
return wantSpec[i].(map[string]interface{})["parameter"].(string) < wantSpec[j].(map[string]interface{})["parameter"].(string)
})
if !reflect.DeepEqual(got.Object["kind"], tt.want.Object["kind"]) &&
!reflect.DeepEqual(got.Object["apiVersion"], tt.want.Object["apiVersion"]) &&
!reflect.DeepEqual(got.Object["metadata"], tt.want.Object["metadata"]) &&
!reflect.DeepEqual(gotSpec, wantSpec) {
t.Errorf("getAuthTriggerSpec() = %v, want %v", got, tt.want)
}
})
}
}
+43
View File
@@ -0,0 +1,43 @@
package util
import (
"fmt"
"os"
"strconv"
"strings"
)
// FissionMetadata contains common fission side fields
type FissionMetadata struct {
// fission
Topic string
ResponseTopic string
ErrorTopic string
FunctionURL string
MaxRetries int
ContentType string
TriggerName string
}
// ParseFissionMetadata parses fission side common fields and returns as fissionMetadata or returns error
func ParseFissionMetadata() (FissionMetadata, error) {
for _, envVars := range []string{"TOPIC", "FUNCTION_URL", "MAX_RETRIES", "CONTENT_TYPE", "TRIGGER_NAME"} {
if os.Getenv(envVars) == "" {
return FissionMetadata{}, fmt.Errorf("environment variable not found: %v", envVars)
}
}
meta := FissionMetadata{
Topic: os.Getenv("TOPIC"),
ResponseTopic: os.Getenv("RESPONSE_TOPIC"),
ErrorTopic: os.Getenv("ERROR_TOPIC"),
FunctionURL: os.Getenv("FUNCTION_URL"),
ContentType: os.Getenv("CONTENT_TYPE"),
TriggerName: os.Getenv("TRIGGER_NAME"),
}
val, err := strconv.ParseInt(strings.TrimSpace(os.Getenv("MAX_RETRIES")), 0, 64)
if err != nil {
return FissionMetadata{}, fmt.Errorf("failed to parse value from MAX_RETRIES environment variable %v", err)
}
meta.MaxRetries = int(val)
return meta, nil
}