Removed code related to mqtrigger (#1680)
This commit is contained in:
@@ -65,7 +65,6 @@ require (
|
|||||||
github.com/stretchr/testify v1.5.1
|
github.com/stretchr/testify v1.5.1
|
||||||
github.com/ulikunitz/xz v0.0.0-20180703112113-636d36a76670 // indirect
|
github.com/ulikunitz/xz v0.0.0-20180703112113-636d36a76670 // indirect
|
||||||
github.com/wcharczuk/go-chart v2.0.1+incompatible
|
github.com/wcharczuk/go-chart v2.0.1+incompatible
|
||||||
github.com/xdg/scram v0.0.0-20180814205039-7eeb5667e42c
|
|
||||||
go.opencensus.io v0.22.0
|
go.opencensus.io v0.22.0
|
||||||
go.uber.org/atomic v1.3.2 // indirect
|
go.uber.org/atomic v1.3.2 // indirect
|
||||||
go.uber.org/multierr v1.1.0 // indirect
|
go.uber.org/multierr v1.1.0 // indirect
|
||||||
@@ -83,4 +82,4 @@ require (
|
|||||||
k8s.io/apimachinery v0.0.0-20190612205821-1799e75a0719
|
k8s.io/apimachinery v0.0.0-20190612205821-1799e75a0719
|
||||||
k8s.io/client-go v12.0.0+incompatible
|
k8s.io/client-go v12.0.0+incompatible
|
||||||
k8s.io/klog v0.3.3
|
k8s.io/klog v0.3.3
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -494,14 +494,14 @@ func (spec MessageQueueTriggerSpec) Validate() error {
|
|||||||
|
|
||||||
result = multierror.Append(result, spec.FunctionReference.Validate())
|
result = multierror.Append(result, spec.FunctionReference.Validate())
|
||||||
|
|
||||||
if !validator.IsValidMessageQueue((string)(spec.MessageQueueType)) {
|
if !validator.IsValidMessageQueue((string)(spec.MessageQueueType), spec.MqtKind) {
|
||||||
result = multierror.Append(result, MakeValidationErr(ErrorUnsupportedType, "MessageQueueTriggerSpec.MessageQueueType", spec.MessageQueueType, "not a supported message queue type"))
|
result = multierror.Append(result, MakeValidationErr(ErrorUnsupportedType, "MessageQueueTriggerSpec.MessageQueueType", spec.MessageQueueType, "not a supported message queue type"))
|
||||||
} else {
|
} else {
|
||||||
if !validator.IsValidTopic((string)(spec.MessageQueueType), spec.Topic) {
|
if !validator.IsValidTopic((string)(spec.MessageQueueType), spec.Topic, spec.MqtKind) {
|
||||||
result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "MessageQueueTriggerSpec.Topic", spec.Topic, "not a valid topic"))
|
result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "MessageQueueTriggerSpec.Topic", spec.Topic, "not a valid topic"))
|
||||||
}
|
}
|
||||||
|
|
||||||
if len(spec.ResponseTopic) > 0 && !validator.IsValidTopic((string)(spec.MessageQueueType), spec.ResponseTopic) {
|
if len(spec.ResponseTopic) > 0 && !validator.IsValidTopic((string)(spec.MessageQueueType), spec.ResponseTopic, spec.MqtKind) {
|
||||||
result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "MessageQueueTriggerSpec.ResponseTopic", spec.ResponseTopic, "not a valid topic"))
|
result = multierror.Append(result, MakeValidationErr(ErrorInvalidValue, "MessageQueueTriggerSpec.ResponseTopic", spec.ResponseTopic, "not a valid topic"))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -59,8 +59,10 @@ func (opts *CreateSubCommand) complete(input cli.Input) error {
|
|||||||
fnName := input.String(flagkey.MqtFnName)
|
fnName := input.String(flagkey.MqtFnName)
|
||||||
fnNamespace := input.String(flagkey.NamespaceFunction)
|
fnNamespace := input.String(flagkey.NamespaceFunction)
|
||||||
|
|
||||||
|
mqtKind := input.String(flagkey.MqtKind)
|
||||||
|
|
||||||
mqType := (fv1.MessageQueueType)(input.String(flagkey.MqtMQType))
|
mqType := (fv1.MessageQueueType)(input.String(flagkey.MqtMQType))
|
||||||
if !validator.IsValidMessageQueue((string)(mqType)) {
|
if !validator.IsValidMessageQueue((string)(mqType), mqtKind) {
|
||||||
return errors.New("Unsupported message queue type")
|
return errors.New("Unsupported message queue type")
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -88,7 +90,7 @@ func (opts *CreateSubCommand) complete(input cli.Input) error {
|
|||||||
contentType = "application/json"
|
contentType = "application/json"
|
||||||
}
|
}
|
||||||
|
|
||||||
err := checkMQTopicAvailability(mqType, topic, respTopic)
|
err := checkMQTopicAvailability(mqType, mqtKind, topic, respTopic)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
@@ -119,8 +121,6 @@ func (opts *CreateSubCommand) complete(input cli.Input) error {
|
|||||||
|
|
||||||
secret := input.String(flagkey.MqtSecret)
|
secret := input.String(flagkey.MqtSecret)
|
||||||
|
|
||||||
mqtKind := input.String(flagkey.MqtKind)
|
|
||||||
|
|
||||||
if input.Bool(flagkey.SpecSave) {
|
if input.Bool(flagkey.SpecSave) {
|
||||||
specDir := util.GetSpecDir(input)
|
specDir := util.GetSpecDir(input)
|
||||||
fr, err := spec.ReadSpecs(specDir)
|
fr, err := spec.ReadSpecs(specDir)
|
||||||
@@ -197,9 +197,9 @@ func (opts *CreateSubCommand) run(input cli.Input) error {
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func checkMQTopicAvailability(mqType fv1.MessageQueueType, topics ...string) error {
|
func checkMQTopicAvailability(mqType fv1.MessageQueueType, mqtKind string, topics ...string) error {
|
||||||
for _, t := range topics {
|
for _, t := range topics {
|
||||||
if len(t) > 0 && !validator.IsValidTopic((string)(mqType), t) {
|
if len(t) > 0 && !validator.IsValidTopic((string)(mqType), t, mqtKind) {
|
||||||
return errors.Errorf("invalid topic for %s: %s", mqType, t)
|
return errors.Errorf("invalid topic for %s: %s", mqType, t)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,19 +0,0 @@
|
|||||||
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"]
|
|
||||||
@@ -1,453 +0,0 @@
|
|||||||
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))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
@@ -54,7 +54,7 @@ func getAuthTriggerClient(namespace string) (dynamic.ResourceInterface, error) {
|
|||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
return dynamicClient.Resource(scaledObjectGVR).Namespace(namespace), nil
|
return dynamicClient.Resource(authTriggerGVR).Namespace(namespace), nil
|
||||||
}
|
}
|
||||||
|
|
||||||
// StartScalerManager watches for changes in MessageQueueTrigger and,
|
// StartScalerManager watches for changes in MessageQueueTrigger and,
|
||||||
@@ -166,7 +166,7 @@ func getEnvVarlist(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient ku
|
|||||||
Value: mqt.Spec.Topic,
|
Value: mqt.Spec.Topic,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
Name: "FUNCTION_URL",
|
Name: "HTTP_ENDPOINT",
|
||||||
Value: url,
|
Value: url,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
@@ -178,7 +178,7 @@ func getEnvVarlist(mqt *fv1.MessageQueueTrigger, routerURL string, kubeClient ku
|
|||||||
Value: mqt.Spec.ResponseTopic,
|
Value: mqt.Spec.ResponseTopic,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
Name: "TRIGGER_NAME",
|
Name: "SOURCE_NAME",
|
||||||
Value: mqt.ObjectMeta.Name,
|
Value: mqt.ObjectMeta.Name,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -111,7 +111,7 @@ func Test_getEnvVarlist(t *testing.T) {
|
|||||||
Value: mqt.Spec.Topic,
|
Value: mqt.Spec.Topic,
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
Name: "FUNCTION_URL",
|
Name: "HTTP_ENDPOINT",
|
||||||
Value: "http://router.fission/fission-function/fission-function/test",
|
Value: "http://router.fission/fission-function/fission-function/test",
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
@@ -123,7 +123,7 @@ func Test_getEnvVarlist(t *testing.T) {
|
|||||||
Value: "response-topic",
|
Value: "response-topic",
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
Name: "TRIGGER_NAME",
|
Name: "SOURCE_NAME",
|
||||||
Value: "Test",
|
Value: "Test",
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
|
|||||||
@@ -1,43 +0,0 @@
|
|||||||
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
|
|
||||||
}
|
|
||||||
@@ -45,7 +45,10 @@ func Register(mqType string, validator TopicValidator) {
|
|||||||
topicValidators[mqType] = validator
|
topicValidators[mqType] = validator
|
||||||
}
|
}
|
||||||
|
|
||||||
func IsValidTopic(mqType string, topic string) bool {
|
func IsValidTopic(mqType, topic, mqtKind string) bool {
|
||||||
|
if mqtKind == "keda" {
|
||||||
|
return true
|
||||||
|
}
|
||||||
validator, registered := topicValidators[mqType]
|
validator, registered := topicValidators[mqType]
|
||||||
if !registered {
|
if !registered {
|
||||||
return false
|
return false
|
||||||
@@ -53,7 +56,10 @@ func IsValidTopic(mqType string, topic string) bool {
|
|||||||
return validator(topic)
|
return validator(topic)
|
||||||
}
|
}
|
||||||
|
|
||||||
func IsValidMessageQueue(mqType string) bool {
|
func IsValidMessageQueue(mqType, mqtKind string) bool {
|
||||||
|
if mqtKind == "keda" {
|
||||||
|
return true
|
||||||
|
}
|
||||||
_, registered := topicValidators[mqType]
|
_, registered := topicValidators[mqType]
|
||||||
return registered
|
return registered
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user