Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20,676 changes: 10,336 additions & 10,340 deletions generate/zz_filesystem_generated.go

Large diffs are not rendered by default.

28 changes: 28 additions & 0 deletions pkg/docker/runner.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
44 changes: 41 additions & 3 deletions pkg/functions/function.go
Original file line number Diff line number Diff line change
Expand Up @@ -180,9 +180,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 }} syntax"`
Password string `yaml:"password,omitempty" jsonschema:"description=SASL password. Supports {{ secret:name:key }} syntax"`
}
Comment on lines +198 to 202

func validateKafka(kafka *KafkaConfig, invoke, runtime string) (errors []string) {
Expand All @@ -206,6 +222,28 @@ 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")
}
}
Comment on lines +237 to +245

return
}

Expand Down
76 changes: 76 additions & 0 deletions pkg/functions/function_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -666,6 +666,82 @@ 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",
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
Expand Down
28 changes: 28 additions & 0 deletions pkg/functions/runner.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
65 changes: 61 additions & 4 deletions pkg/k8s/deployer.go
Original file line number Diff line number Diff line change
Expand Up @@ -428,7 +428,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 {
Expand Down Expand Up @@ -732,17 +735,71 @@ 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"},
corev1.EnvVar{Name: "KAFKA_BROKERS", Value: kafka.Brokers},
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, "{{") {
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 %q, expected {{ secret:name:key }} or {{ configMap:name:key }}", value)
}
return append(envVars, corev1.EnvVar{Name: name, Value: value}), nil
}
Comment on lines +790 to 803

func createEnvFromSource(value string, referencedSecrets, referencedConfigMaps *sets.Set[string]) (*corev1.EnvFromSource, error) {
Expand Down
Loading
Loading