diff --git a/pkg/docker/runner.go b/pkg/docker/runner.go index 3a2c5b4e82..c2bc3fff22 100644 --- a/pkg/docker/runner.go +++ b/pkg/docker/runner.go @@ -306,6 +306,34 @@ func newContainerConfig(f fn.Function, _ string, verbose bool) (c container.Conf "KAFKA_TOPIC="+k.Topic, "KAFKA_CONSUMER_GROUP="+k.ConsumerGroup, ) + if k.SecurityProtocol != "" && k.SecurityProtocol != "PLAINTEXT" { + c.Env = append(c.Env, "KAFKA_SECURITY_PROTOCOL="+k.SecurityProtocol) + } + if k.TLS != nil { + if k.TLS.CACert != "" { + c.Env = append(c.Env, "KAFKA_TLS_CA_CERT="+k.TLS.CACert) + } + if k.TLS.ClientCert != "" { + c.Env = append(c.Env, "KAFKA_TLS_CLIENT_CERT="+k.TLS.ClientCert) + } + if k.TLS.ClientKey != "" { + c.Env = append(c.Env, "KAFKA_TLS_CLIENT_KEY="+k.TLS.ClientKey) + } + if k.TLS.SkipVerify { + c.Env = append(c.Env, "KAFKA_TLS_SKIP_VERIFY=true") + } + } + if k.SASL != nil { + if k.SASL.Mechanism != "" { + c.Env = append(c.Env, "KAFKA_SASL_MECHANISM="+k.SASL.Mechanism) + } + if k.SASL.User != "" { + c.Env = append(c.Env, "KAFKA_SASL_USER="+k.SASL.User) + } + if k.SASL.Password != "" { + c.Env = append(c.Env, "KAFKA_SASL_PASSWORD="+k.SASL.Password) + } + } } return diff --git a/pkg/functions/function.go b/pkg/functions/function.go index 02bd45d2e4..de2833be98 100644 --- a/pkg/functions/function.go +++ b/pkg/functions/function.go @@ -186,9 +186,25 @@ type MountSpec struct { // When set, the runtime consumes messages from Kafka and delivers them // as CloudEvents to the function's handler. type KafkaConfig struct { - Brokers string `yaml:"brokers" jsonschema:"description=Comma-separated list of Kafka broker addresses"` - Topic string `yaml:"topic" jsonschema:"description=Kafka topic to consume from"` - ConsumerGroup string `yaml:"consumerGroup" jsonschema:"description=Kafka consumer group ID"` + Brokers string `yaml:"brokers" jsonschema:"description=Comma-separated list of Kafka broker addresses"` + Topic string `yaml:"topic" jsonschema:"description=Kafka topic to consume from"` + ConsumerGroup string `yaml:"consumerGroup" jsonschema:"description=Kafka consumer group ID"` + SecurityProtocol string `yaml:"securityProtocol,omitempty" jsonschema:"description=Security protocol: PLAINTEXT SSL SASL_PLAINTEXT or SASL_SSL,enum=PLAINTEXT,enum=SSL,enum=SASL_PLAINTEXT,enum=SASL_SSL"` + TLS *KafkaTLS `yaml:"tls,omitempty" jsonschema:"description=TLS configuration for SSL or SASL_SSL"` + SASL *KafkaSASL `yaml:"sasl,omitempty" jsonschema:"description=SASL authentication for SASL_PLAINTEXT or SASL_SSL"` +} + +type KafkaTLS struct { + CACert string `yaml:"caCert,omitempty" jsonschema:"description=Path to CA certificate PEM file for verifying broker certificate"` + ClientCert string `yaml:"clientCert,omitempty" jsonschema:"description=Path to client certificate PEM file for mutual TLS"` + ClientKey string `yaml:"clientKey,omitempty" jsonschema:"description=Path to client private key PEM file for mutual TLS"` + SkipVerify bool `yaml:"skipVerify,omitempty" jsonschema:"description=Skip broker certificate verification (development only)"` +} + +type KafkaSASL struct { + Mechanism string `yaml:"mechanism,omitempty" jsonschema:"description=SASL mechanism: PLAIN SCRAM-SHA-256 or SCRAM-SHA-512,enum=PLAIN,enum=SCRAM-SHA-256,enum=SCRAM-SHA-512"` + User string `yaml:"user,omitempty" jsonschema:"description=SASL username. Supports {{ secret:name:key }} and {{ configMap:name:key }} syntax"` + Password string `yaml:"password,omitempty" jsonschema:"description=SASL password. Supports {{ secret:name:key }} and {{ configMap:name:key }} syntax"` } func validateKafka(kafka *KafkaConfig, invoke, runtime string) (errors []string) { @@ -212,6 +228,39 @@ func validateKafka(kafka *KafkaConfig, invoke, runtime string) (errors []string) if kafka.ConsumerGroup == "" { errors = append(errors, "run.kafka.consumerGroup is required when Kafka is configured") } + + validProtocols := map[string]bool{"": true, "PLAINTEXT": true, "SSL": true, "SASL_PLAINTEXT": true, "SASL_SSL": true} + if !validProtocols[kafka.SecurityProtocol] { + errors = append(errors, "run.kafka.securityProtocol must be one of: PLAINTEXT, SSL, SASL_PLAINTEXT, SASL_SSL") + } + + if kafka.TLS != nil { + if kafka.SecurityProtocol != "SSL" && kafka.SecurityProtocol != "SASL_SSL" { + errors = append(errors, "run.kafka.tls requires securityProtocol SSL or SASL_SSL") + } + } + + if kafka.SASL != nil { + if kafka.SecurityProtocol != "SASL_PLAINTEXT" && kafka.SecurityProtocol != "SASL_SSL" { + errors = append(errors, "run.kafka.sasl requires securityProtocol SASL_PLAINTEXT or SASL_SSL") + } + validMechanisms := map[string]bool{"": true, "PLAIN": true, "SCRAM-SHA-256": true, "SCRAM-SHA-512": true} + if !validMechanisms[kafka.SASL.Mechanism] { + errors = append(errors, "run.kafka.sasl.mechanism must be one of: PLAIN, SCRAM-SHA-256, SCRAM-SHA-512") + } + errors = append(errors, validateTemplateRef("run.kafka.sasl.user", kafka.SASL.User)...) + errors = append(errors, validateTemplateRef("run.kafka.sasl.password", kafka.SASL.Password)...) + } + + return +} + +var templateRefPattern = regexp.MustCompile(`^\{\{\s*(secret|configMap):[^:]+:[^:]+\s*\}\}$`) + +func validateTemplateRef(field, value string) (errors []string) { + if strings.HasPrefix(value, "{{") && !templateRefPattern.MatchString(value) { + errors = append(errors, fmt.Sprintf("%s has invalid reference format, expected {{ secret:name:key }} or {{ configMap:name:key }}", field)) + } return } diff --git a/pkg/functions/function_test.go b/pkg/functions/function_test.go index 3a85a730af..9393d99dfe 100644 --- a/pkg/functions/function_test.go +++ b/pkg/functions/function_test.go @@ -666,6 +666,108 @@ func TestValidateKafka(t *testing.T) { wantErrs: 3, wantSubst: "required", }, + { + name: "valid SASL_SSL config", + kafka: &fn.KafkaConfig{ + Brokers: "broker:9093", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "SASL_SSL", + TLS: &fn.KafkaTLS{CACert: "/etc/kafka/ca/ca.crt"}, + SASL: &fn.KafkaSASL{Mechanism: "SCRAM-SHA-512", User: "u", Password: "p"}, + }, + invoke: "cloudevent", + wantErrs: 0, + }, + { + name: "valid SSL config", + kafka: &fn.KafkaConfig{ + Brokers: "broker:9093", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "SSL", + TLS: &fn.KafkaTLS{CACert: "/etc/kafka/ca/ca.crt"}, + }, + invoke: "cloudevent", + wantErrs: 0, + }, + { + name: "invalid security protocol", + kafka: &fn.KafkaConfig{ + Brokers: "broker:9092", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "BOGUS", + }, + invoke: "cloudevent", + wantErrs: 1, + wantSubst: "securityProtocol must be one of", + }, + { + name: "tls without SSL protocol", + kafka: &fn.KafkaConfig{ + Brokers: "broker:9092", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "PLAINTEXT", + TLS: &fn.KafkaTLS{CACert: "/ca.crt"}, + }, + invoke: "cloudevent", + wantErrs: 1, + wantSubst: "run.kafka.tls requires securityProtocol SSL or SASL_SSL", + }, + { + name: "sasl without SASL protocol", + kafka: &fn.KafkaConfig{ + Brokers: "broker:9092", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "SSL", + SASL: &fn.KafkaSASL{Mechanism: "PLAIN", User: "u", Password: "p"}, + }, + invoke: "cloudevent", + wantErrs: 1, + wantSubst: "run.kafka.sasl requires securityProtocol SASL_PLAINTEXT or SASL_SSL", + }, + { + name: "invalid SASL mechanism", + kafka: &fn.KafkaConfig{ + Brokers: "broker:9092", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "SASL_SSL", + SASL: &fn.KafkaSASL{Mechanism: "OAUTHBEARER"}, + }, + invoke: "cloudevent", + wantErrs: 1, + wantSubst: "sasl.mechanism must be one of", + }, + { + name: "malformed secret ref in SASL password", + kafka: &fn.KafkaConfig{ + Brokers: "broker:9093", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "SASL_SSL", + SASL: &fn.KafkaSASL{Mechanism: "PLAIN", User: "u", Password: "{{ secret:s:k }}extra"}, + }, + invoke: "cloudevent", + wantErrs: 1, + wantSubst: "invalid reference format", + }, + { + name: "valid secret ref in SASL password", + kafka: &fn.KafkaConfig{ + Brokers: "broker:9093", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "SASL_SSL", + TLS: &fn.KafkaTLS{CACert: "/ca.crt"}, + SASL: &fn.KafkaSASL{Mechanism: "PLAIN", User: "u", Password: "{{ secret:my-user:password }}"}, + }, + invoke: "cloudevent", + wantErrs: 0, + }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { diff --git a/pkg/functions/runner.go b/pkg/functions/runner.go index 6b7286308d..04091b70cd 100644 --- a/pkg/functions/runner.go +++ b/pkg/functions/runner.go @@ -314,6 +314,34 @@ func buildRunnerEnv(job *Job, extras map[string]string) ([]string, error) { "KAFKA_TOPIC="+k.Topic, "KAFKA_CONSUMER_GROUP="+k.ConsumerGroup, ) + if k.SecurityProtocol != "" && k.SecurityProtocol != "PLAINTEXT" { + env = append(env, "KAFKA_SECURITY_PROTOCOL="+k.SecurityProtocol) + } + if k.TLS != nil { + if k.TLS.CACert != "" { + env = append(env, "KAFKA_TLS_CA_CERT="+k.TLS.CACert) + } + if k.TLS.ClientCert != "" { + env = append(env, "KAFKA_TLS_CLIENT_CERT="+k.TLS.ClientCert) + } + if k.TLS.ClientKey != "" { + env = append(env, "KAFKA_TLS_CLIENT_KEY="+k.TLS.ClientKey) + } + if k.TLS.SkipVerify { + env = append(env, "KAFKA_TLS_SKIP_VERIFY=true") + } + } + if k.SASL != nil { + if k.SASL.Mechanism != "" { + env = append(env, "KAFKA_SASL_MECHANISM="+k.SASL.Mechanism) + } + if k.SASL.User != "" { + env = append(env, "KAFKA_SASL_USER="+k.SASL.User) + } + if k.SASL.Password != "" { + env = append(env, "KAFKA_SASL_PASSWORD="+k.SASL.Password) + } + } } return env, nil diff --git a/pkg/k8s/deployer.go b/pkg/k8s/deployer.go index b7102557a7..32740b9283 100644 --- a/pkg/k8s/deployer.go +++ b/pkg/k8s/deployer.go @@ -430,7 +430,10 @@ func (d *Deployer) generateDeployment(f fn.Function, namespace string, daprInsta if err != nil { return nil, fmt.Errorf("failed to process environment variables: %w", err) } - envVars = AppendKafkaEnvs(envVars, f.Run.Kafka) + envVars, err = AppendKafkaEnvs(envVars, f.Run.Kafka, referencedSecrets, referencedConfigMaps) + if err != nil { + return nil, fmt.Errorf("failed to process Kafka environment variables: %w", err) + } volumes, volumeMounts, err := ProcessVolumes(f.Run.Volumes, referencedSecrets, referencedConfigMaps, referencedPVCs) if err != nil { @@ -734,9 +737,9 @@ func withOpenAddress(ee []fn.Env) []fn.Env { return ee } -func AppendKafkaEnvs(envVars []corev1.EnvVar, kafka *fn.KafkaConfig) []corev1.EnvVar { +func AppendKafkaEnvs(envVars []corev1.EnvVar, kafka *fn.KafkaConfig, referencedSecrets, referencedConfigMaps *sets.Set[string]) ([]corev1.EnvVar, error) { if kafka == nil || kafka.Brokers == "" || kafka.Topic == "" || kafka.ConsumerGroup == "" { - return envVars + return envVars, nil } envVars = append(envVars, corev1.EnvVar{Name: "FUNC_TRANSPORT", Value: "kafka"}, @@ -744,7 +747,64 @@ func AppendKafkaEnvs(envVars []corev1.EnvVar, kafka *fn.KafkaConfig) []corev1.En corev1.EnvVar{Name: "KAFKA_TOPIC", Value: kafka.Topic}, corev1.EnvVar{Name: "KAFKA_CONSUMER_GROUP", Value: kafka.ConsumerGroup}, ) - return envVars + + if kafka.SecurityProtocol != "" && kafka.SecurityProtocol != "PLAINTEXT" { + envVars = append(envVars, corev1.EnvVar{Name: "KAFKA_SECURITY_PROTOCOL", Value: kafka.SecurityProtocol}) + } + + if kafka.TLS != nil { + if kafka.TLS.CACert != "" { + envVars = append(envVars, corev1.EnvVar{Name: "KAFKA_TLS_CA_CERT", Value: kafka.TLS.CACert}) + } + if kafka.TLS.ClientCert != "" { + envVars = append(envVars, corev1.EnvVar{Name: "KAFKA_TLS_CLIENT_CERT", Value: kafka.TLS.ClientCert}) + } + if kafka.TLS.ClientKey != "" { + envVars = append(envVars, corev1.EnvVar{Name: "KAFKA_TLS_CLIENT_KEY", Value: kafka.TLS.ClientKey}) + } + if kafka.TLS.SkipVerify { + envVars = append(envVars, corev1.EnvVar{Name: "KAFKA_TLS_SKIP_VERIFY", Value: "true"}) + } + } + + if kafka.SASL != nil { + if kafka.SASL.Mechanism != "" { + envVars = append(envVars, corev1.EnvVar{Name: "KAFKA_SASL_MECHANISM", Value: kafka.SASL.Mechanism}) + } + var err error + if kafka.SASL.User != "" { + envVars, err = appendKafkaEnvValue(envVars, "KAFKA_SASL_USER", kafka.SASL.User, referencedSecrets, referencedConfigMaps) + if err != nil { + return nil, fmt.Errorf("processing run.kafka.sasl.user: %w", err) + } + } + if kafka.SASL.Password != "" { + envVars, err = appendKafkaEnvValue(envVars, "KAFKA_SASL_PASSWORD", kafka.SASL.Password, referencedSecrets, referencedConfigMaps) + if err != nil { + return nil, fmt.Errorf("processing run.kafka.sasl.password: %w", err) + } + } + } + + return envVars, nil +} + +func appendKafkaEnvValue(envVars []corev1.EnvVar, name, value string, referencedSecrets, referencedConfigMaps *sets.Set[string]) ([]corev1.EnvVar, error) { + if strings.HasPrefix(value, "{{") { + if !strings.HasSuffix(strings.TrimSpace(value), "}}") { + return nil, fmt.Errorf("invalid reference format for %s, expected {{ secret:name:key }} or {{ configMap:name:key }}", name) + } + slices := strings.Split(strings.Trim(value, "{} "), ":") + if len(slices) == 3 { + valueFrom, err := createEnvVarSource(slices, referencedSecrets, referencedConfigMaps) + if err != nil { + return nil, err + } + return append(envVars, corev1.EnvVar{Name: name, ValueFrom: valueFrom}), nil + } + return nil, fmt.Errorf("invalid reference format for %s, expected {{ secret:name:key }} or {{ configMap:name:key }}", name) + } + return append(envVars, corev1.EnvVar{Name: name, Value: value}), nil } func createEnvFromSource(value string, referencedSecrets, referencedConfigMaps *sets.Set[string]) (*corev1.EnvFromSource, error) { diff --git a/pkg/k8s/deployer_test.go b/pkg/k8s/deployer_test.go index 020abc2d9d..7dccaef71a 100644 --- a/pkg/k8s/deployer_test.go +++ b/pkg/k8s/deployer_test.go @@ -430,7 +430,12 @@ func TestAppendKafkaEnvs_Nil(t *testing.T) { base := []corev1.EnvVar{ {Name: "EXISTING", Value: "value"}, } - got := AppendKafkaEnvs(base, nil) + secrets := sets.New[string]() + configMaps := sets.New[string]() + got, err := AppendKafkaEnvs(base, nil, &secrets, &configMaps) + if err != nil { + t.Fatal(err) + } if len(got) != 1 { t.Fatalf("expected 1 env var, got %d", len(got)) } @@ -448,7 +453,12 @@ func TestAppendKafkaEnvs_AllFields(t *testing.T) { Topic: "my-topic", ConsumerGroup: "my-group", } - got := AppendKafkaEnvs(base, kafka) + secrets := sets.New[string]() + configMaps := sets.New[string]() + got, err := AppendKafkaEnvs(base, kafka, &secrets, &configMaps) + if err != nil { + t.Fatal(err) + } if len(got) != 5 { t.Fatalf("expected 5 env vars (1 existing + 4 kafka), got %d", len(got)) } @@ -479,7 +489,12 @@ func TestAppendKafkaEnvs_MissingBrokers(t *testing.T) { Topic: "my-topic", ConsumerGroup: "my-group", } - got := AppendKafkaEnvs(base, kafka) + secrets := sets.New[string]() + configMaps := sets.New[string]() + got, err := AppendKafkaEnvs(base, kafka, &secrets, &configMaps) + if err != nil { + t.Fatal(err) + } if len(got) != 1 { t.Fatalf("expected 1 env var (unchanged), got %d", len(got)) } @@ -494,12 +509,118 @@ func TestAppendKafkaEnvs_MissingTopic(t *testing.T) { Topic: "", ConsumerGroup: "my-group", } - got := AppendKafkaEnvs(base, kafka) + secrets := sets.New[string]() + configMaps := sets.New[string]() + got, err := AppendKafkaEnvs(base, kafka, &secrets, &configMaps) + if err != nil { + t.Fatal(err) + } if len(got) != 1 { t.Fatalf("expected 1 env var (unchanged), got %d", len(got)) } } +func TestAppendKafkaEnvs_SASL_SSL(t *testing.T) { + kafka := &fn.KafkaConfig{ + Brokers: "broker:9093", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "SASL_SSL", + TLS: &fn.KafkaTLS{CACert: "/etc/kafka/ca/ca.crt"}, + SASL: &fn.KafkaSASL{Mechanism: "SCRAM-SHA-512", User: "alice", Password: "s3cret"}, + } + secrets := sets.New[string]() + configMaps := sets.New[string]() + got, err := AppendKafkaEnvs(nil, kafka, &secrets, &configMaps) + if err != nil { + t.Fatal(err) + } + envMap := make(map[string]string) + for _, ev := range got { + envMap[ev.Name] = ev.Value + } + if envMap["KAFKA_SECURITY_PROTOCOL"] != "SASL_SSL" { + t.Errorf("KAFKA_SECURITY_PROTOCOL = %q", envMap["KAFKA_SECURITY_PROTOCOL"]) + } + if envMap["KAFKA_TLS_CA_CERT"] != "/etc/kafka/ca/ca.crt" { + t.Errorf("KAFKA_TLS_CA_CERT = %q", envMap["KAFKA_TLS_CA_CERT"]) + } + if envMap["KAFKA_SASL_MECHANISM"] != "SCRAM-SHA-512" { + t.Errorf("KAFKA_SASL_MECHANISM = %q", envMap["KAFKA_SASL_MECHANISM"]) + } + if envMap["KAFKA_SASL_USER"] != "alice" { + t.Errorf("KAFKA_SASL_USER = %q", envMap["KAFKA_SASL_USER"]) + } + if envMap["KAFKA_SASL_PASSWORD"] != "s3cret" { + t.Errorf("KAFKA_SASL_PASSWORD = %q", envMap["KAFKA_SASL_PASSWORD"]) + } +} + +func TestAppendKafkaEnvs_SecretRef(t *testing.T) { + kafka := &fn.KafkaConfig{ + Brokers: "broker:9093", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "SASL_SSL", + SASL: &fn.KafkaSASL{ + Mechanism: "PLAIN", + User: "{{ secret:my-kafka-user:username }}", + Password: "{{ secret:my-kafka-user:password }}", + }, + } + secrets := sets.New[string]() + configMaps := sets.New[string]() + got, err := AppendKafkaEnvs(nil, kafka, &secrets, &configMaps) + if err != nil { + t.Fatal(err) + } + + for _, ev := range got { + if ev.Name == "KAFKA_SASL_USER" { + if ev.ValueFrom == nil || ev.ValueFrom.SecretKeyRef == nil { + t.Fatal("KAFKA_SASL_USER should have ValueFrom with SecretKeyRef") + } + if ev.ValueFrom.SecretKeyRef.Name != "my-kafka-user" { + t.Errorf("secret name = %q, want my-kafka-user", ev.ValueFrom.SecretKeyRef.Name) + } + if ev.ValueFrom.SecretKeyRef.Key != "username" { + t.Errorf("secret key = %q, want username", ev.ValueFrom.SecretKeyRef.Key) + } + } + if ev.Name == "KAFKA_SASL_PASSWORD" { + if ev.ValueFrom == nil || ev.ValueFrom.SecretKeyRef == nil { + t.Fatal("KAFKA_SASL_PASSWORD should have ValueFrom with SecretKeyRef") + } + if ev.ValueFrom.SecretKeyRef.Key != "password" { + t.Errorf("secret key = %q, want password", ev.ValueFrom.SecretKeyRef.Key) + } + } + } + if !secrets.Has("my-kafka-user") { + t.Error("secret my-kafka-user should be tracked in referencedSecrets") + } +} + +func TestAppendKafkaEnvs_MalformedSecretRef(t *testing.T) { + kafka := &fn.KafkaConfig{ + Brokers: "broker:9093", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "SASL_SSL", + SASL: &fn.KafkaSASL{ + Mechanism: "PLAIN", + User: "{{ secret:my-user:name }}trailing", + Password: "p", + }, + } + secrets := sets.New[string]() + configMaps := sets.New[string]() + _, err := AppendKafkaEnvs(nil, kafka, &secrets, &configMaps) + if err == nil { + t.Fatal("expected error for malformed secret reference with trailing content") + } +} + func Test_ProcessVolumes_NilPath(t *testing.T) { secretName := "my-secret" referencedSecrets := sets.New[string]() diff --git a/pkg/knative/deployer.go b/pkg/knative/deployer.go index 7d47065873..879196f98b 100644 --- a/pkg/knative/deployer.go +++ b/pkg/knative/deployer.go @@ -315,7 +315,10 @@ consider using the --image-pull-secret flag, or setting up pull secrets manually if err != nil { return fn.DeploymentResult{}, err } - newEnv = k8s.AppendKafkaEnvs(newEnv, f.Run.Kafka) + newEnv, err = k8s.AppendKafkaEnvs(newEnv, f.Run.Kafka, &referencedSecrets, &referencedConfigMaps) + if err != nil { + return fn.DeploymentResult{}, err + } newVolumes, newVolumeMounts, err := k8s.ProcessVolumes(f.Run.Volumes, &referencedSecrets, &referencedConfigMaps, &referencedPVCs) if err != nil { @@ -429,7 +432,10 @@ func generateNewService(f fn.Function, decorator deployer.DeployDecorator, daprI if err != nil { return nil, err } - container.Env = k8s.AppendKafkaEnvs(newEnv, f.Run.Kafka) + container.Env, err = k8s.AppendKafkaEnvs(newEnv, f.Run.Kafka, referencedSecrets, referencedConfigMaps) + if err != nil { + return nil, err + } container.EnvFrom = newEnvFrom newVolumes, newVolumeMounts, err := k8s.ProcessVolumes(f.Run.Volumes, referencedSecrets, referencedConfigMaps, referencedPVCs) diff --git a/schema/func_yaml-schema.json b/schema/func_yaml-schema.json index 044a19b699..f2dff0559d 100644 --- a/schema/func_yaml-schema.json +++ b/schema/func_yaml-schema.json @@ -290,12 +290,77 @@ "consumerGroup": { "type": "string", "description": "Kafka consumer group ID" + }, + "securityProtocol": { + "enum": [ + "PLAINTEXT", + "SSL", + "SASL_PLAINTEXT", + "SASL_SSL" + ], + "type": "string", + "description": "Security protocol: PLAINTEXT SSL SASL_PLAINTEXT or SASL_SSL" + }, + "tls": { + "$schema": "http://json-schema.org/draft-04/schema#", + "$ref": "#/definitions/KafkaTLS", + "description": "TLS configuration for SSL or SASL_SSL" + }, + "sasl": { + "$schema": "http://json-schema.org/draft-04/schema#", + "$ref": "#/definitions/KafkaSASL", + "description": "SASL authentication for SASL_PLAINTEXT or SASL_SSL" } }, "additionalProperties": false, "type": "object", "description": "KafkaConfig specifies the Kafka event source configuration." }, + "KafkaSASL": { + "properties": { + "mechanism": { + "enum": [ + "PLAIN", + "SCRAM-SHA-256", + "SCRAM-SHA-512" + ], + "type": "string", + "description": "SASL mechanism: PLAIN SCRAM-SHA-256 or SCRAM-SHA-512" + }, + "user": { + "type": "string", + "description": "SASL username. Supports {{ secret:name:key }} and {{ configMap:name:key }} syntax" + }, + "password": { + "type": "string", + "description": "SASL password. Supports {{ secret:name:key }} and {{ configMap:name:key }} syntax" + } + }, + "additionalProperties": false, + "type": "object" + }, + "KafkaTLS": { + "properties": { + "caCert": { + "type": "string", + "description": "Path to CA certificate PEM file for verifying broker certificate" + }, + "clientCert": { + "type": "string", + "description": "Path to client certificate PEM file for mutual TLS" + }, + "clientKey": { + "type": "string", + "description": "Path to client private key PEM file for mutual TLS" + }, + "skipVerify": { + "type": "boolean", + "description": "Skip broker certificate verification (development only)" + } + }, + "additionalProperties": false, + "type": "object" + }, "KnativeSubscription": { "required": [ "source"