diff --git a/docs/reference/func_yaml.md b/docs/reference/func_yaml.md index 09a9075561..c18991b198 100644 --- a/docs/reference/func_yaml.md +++ b/docs/reference/func_yaml.md @@ -258,7 +258,7 @@ topic instead of serving HTTP requests. Requires `invoke: cloudevent` and the Go - `clientCert`, `clientKey`: paths to the client certificate/key PEM files, for mutual TLS. - `skipVerify`: skip broker certificate verification (development only). - `sasl`: SASL configuration, required for `SASL_PLAINTEXT` and `SASL_SSL`. - - `mechanism`: one of `PLAIN`, `SCRAM-SHA-256`, `SCRAM-SHA-512`. + - `mechanism`: one of `PLAIN`, `SCRAM-SHA-256`, `SCRAM-SHA-512`. Optional; defaults to `PLAIN` when unset. - `user`: SASL username. Supports `{{ secret:name:key }}` and `{{ configMap:name:key }}` syntax, or a plain value. - `password`: SASL password. Supports `{{ secret:name:key }}` and `{{ configMap:name:key }}` syntax, or a plain value (at least for debugging purposes). diff --git a/pkg/docker/runner.go b/pkg/docker/runner.go index 021911630b..1589d1ac0f 100644 --- a/pkg/docker/runner.go +++ b/pkg/docker/runner.go @@ -326,9 +326,7 @@ func newContainerConfig(f fn.Function, _ string, verbose bool) (c container.Conf } } if k.SASL != nil { - if k.SASL.Mechanism != "" { - c.Env = append(c.Env, "KAFKA_SASL_MECHANISM="+k.SASL.Mechanism) - } + c.Env = append(c.Env, "KAFKA_SASL_MECHANISM="+k.SASL.EffectiveMechanism()) if k.SASL.User != "" { c.Env = append(c.Env, "KAFKA_SASL_USER="+k.SASL.User) } diff --git a/pkg/functions/function.go b/pkg/functions/function.go index 8bf8df31d4..d24b01ccbc 100644 --- a/pkg/functions/function.go +++ b/pkg/functions/function.go @@ -214,11 +214,22 @@ type KafkaTLS struct { } 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"` + Mechanism string `yaml:"mechanism,omitempty" jsonschema:"description=SASL mechanism: PLAIN SCRAM-SHA-256 or SCRAM-SHA-512. Optional; defaults to PLAIN when unset.,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"` } +// EffectiveMechanism returns the SASL mechanism the function actually +// authenticates with, resolving an unset mechanism to PLAIN. Callers that wire +// the function's container env and the keda scaler both read it so the two agree +// on a single value rather than each re-deriving the default. +func (s KafkaSASL) EffectiveMechanism() string { + if s.Mechanism == "" { + return "PLAIN" + } + return s.Mechanism +} + func validateKafka(kafka *KafkaConfig, invoke, runtime string) (errors []string) { if kafka == nil { return @@ -276,11 +287,11 @@ func ValidateKafkaSecurity(kafka *KafkaConfig) (errors []string) { if kafka.SecurityProtocol != "SASL_PLAINTEXT" && kafka.SecurityProtocol != "SASL_SSL" { errors = append(errors, "run.kafka.sasl requires securityProtocol SASL_PLAINTEXT or SASL_SSL") } - // mechanism is required: KEDA's TriggerAuthentication and the runtime - // both need a concrete SASL mechanism, there is no sensible default. - if kafka.SASL.Mechanism == "" { - errors = append(errors, "run.kafka.sasl.mechanism is required") - } else { + // mechanism is optional: EffectiveMechanism resolves an unset mechanism + // to PLAIN for both the function container and the scaler, so the two + // authenticate the same way with nothing set. Only a non-empty value is + // constrained to the mechanisms both sides understand. + if kafka.SASL.Mechanism != "" { validMechanisms := map[string]bool{"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") diff --git a/pkg/functions/function_test.go b/pkg/functions/function_test.go index f2be9b1a3f..a21009015a 100644 --- a/pkg/functions/function_test.go +++ b/pkg/functions/function_test.go @@ -743,7 +743,11 @@ func TestValidateKafka(t *testing.T) { wantSubst: "sasl.mechanism must be one of", }, { - name: "empty SASL mechanism is required, not silently accepted", + // An empty mechanism is accepted: EffectiveMechanism resolves it to + // PLAIN for both the container env and the keda scaler, so the two + // agree with nothing set. This matches a plain SASL/PLAIN broker that + // a raw or knative deploy consumed from before scale.keda landed. + name: "empty SASL mechanism defaults to PLAIN", kafka: &fn.KafkaConfig{ Brokers: "broker:9092", Topic: "my-topic", @@ -751,9 +755,8 @@ func TestValidateKafka(t *testing.T) { SecurityProtocol: "SASL_SSL", SASL: &fn.KafkaSASL{User: "u", Password: "p"}, }, - invoke: "cloudevent", - wantErrs: 1, - wantSubst: "run.kafka.sasl.mechanism is required", + invoke: "cloudevent", + wantErrs: 0, }, { name: "TLS clientCert without clientKey", @@ -891,6 +894,25 @@ func TestValidateKafka(t *testing.T) { } } +// TestKafkaSASL_EffectiveMechanism verifies an unset mechanism resolves to +// PLAIN while an explicit value passes through unchanged, so the container env +// and the keda scaler read one agreed value rather than each re-deriving the +// default. +func TestKafkaSASL_EffectiveMechanism(t *testing.T) { + tests := map[string]string{ + "": "PLAIN", + "PLAIN": "PLAIN", + "SCRAM-SHA-256": "SCRAM-SHA-256", + "SCRAM-SHA-512": "SCRAM-SHA-512", + } + for in, want := range tests { + s := fn.KafkaSASL{Mechanism: in} + if got := s.EffectiveMechanism(); got != want { + t.Errorf("EffectiveMechanism() with Mechanism=%q = %q, want %q", in, got, want) + } + } +} + func TestKafkaConfig_YAMLOmitEmpty(t *testing.T) { f := fn.Function{ Name: "test-func", diff --git a/pkg/functions/runner.go b/pkg/functions/runner.go index 41d5ddf569..c01d2790af 100644 --- a/pkg/functions/runner.go +++ b/pkg/functions/runner.go @@ -334,9 +334,7 @@ func buildRunnerEnv(job *Job, extras map[string]string) ([]string, error) { } } if k.SASL != nil { - if k.SASL.Mechanism != "" { - env = append(env, "KAFKA_SASL_MECHANISM="+k.SASL.Mechanism) - } + env = append(env, "KAFKA_SASL_MECHANISM="+k.SASL.EffectiveMechanism()) if k.SASL.User != "" { env = append(env, "KAFKA_SASL_USER="+k.SASL.User) } diff --git a/pkg/functions/runner_test.go b/pkg/functions/runner_test.go index 156c1ecd4f..b36df069ac 100644 --- a/pkg/functions/runner_test.go +++ b/pkg/functions/runner_test.go @@ -207,3 +207,39 @@ func TestBuildRunnerEnv_KafkaTLSSkipVerifyFalse(t *testing.T) { t.Error("expected KAFKA_TLS_SKIP_VERIFY=false when TLS block present with SkipVerify unset") } } + +// TestBuildRunnerEnv_KafkaEmptyMechanism verifies an unset SASL mechanism still +// wires KAFKA_SASL_MECHANISM=PLAIN into the container env, matching the value the +// keda scaler derives from EffectiveMechanism so the two authenticate the same +// way. +func TestBuildRunnerEnv_KafkaEmptyMechanism(t *testing.T) { + job := &Job{ + Function: Function{ + Root: t.TempDir(), + Runtime: "go", + Run: RunSpec{ + Kafka: &KafkaConfig{ + Brokers: "broker:9093", + Topic: "my-topic", + ConsumerGroup: "my-group", + SecurityProtocol: "SASL_SSL", + SASL: &KafkaSASL{User: "u", Password: "p"}, + }, + }, + }, + } + env, err := buildRunnerEnv(job, nil) + if err != nil { + t.Fatal(err) + } + found := false + for _, e := range env { + if e == "KAFKA_SASL_MECHANISM=PLAIN" { + found = true + break + } + } + if !found { + t.Error("expected KAFKA_SASL_MECHANISM=PLAIN when SASL block present with mechanism unset") + } +} diff --git a/pkg/k8s/deployer.go b/pkg/k8s/deployer.go index 12a88cd58b..b98279ba11 100644 --- a/pkg/k8s/deployer.go +++ b/pkg/k8s/deployer.go @@ -1042,9 +1042,7 @@ func AppendKafkaEnvs(envVars []corev1.EnvVar, kafka *fn.KafkaConfig, referencedS } if kafka.SASL != nil { - if kafka.SASL.Mechanism != "" { - envVars = append(envVars, corev1.EnvVar{Name: "KAFKA_SASL_MECHANISM", Value: kafka.SASL.Mechanism}) - } + envVars = append(envVars, corev1.EnvVar{Name: "KAFKA_SASL_MECHANISM", Value: kafka.SASL.EffectiveMechanism()}) var err error if kafka.SASL.User != "" { envVars, err = appendKafkaEnvValue(envVars, "KAFKA_SASL_USER", kafka.SASL.User, referencedSecrets, referencedConfigMaps) diff --git a/pkg/keda/deployer_unit_test.go b/pkg/keda/deployer_unit_test.go index c804e068c1..aea9232c5a 100644 --- a/pkg/keda/deployer_unit_test.go +++ b/pkg/keda/deployer_unit_test.go @@ -352,10 +352,10 @@ func TestPollingIntervalIgnored(t *testing.T) { // TestDeploy_KafkaSASLPreflight covers the Deploy preflight for direct callers // that bypass Function.Validate: an inconsistent SASL/security config must be // rejected before any cluster resources are created. Otherwise buildScaledObject -// -- which only emits "sasl" trigger metadata for a non-empty mechanism -- would -// produce a ScaledObject that connects without SASL while the function's own -// container is configured for it. Each case returns from the pure preflight -// before Deploy touches the cluster, so no fake clientset is needed. +// would produce a ScaledObject whose SASL/TLS trigger metadata silently +// disagrees with how the function's own container authenticates. Each case +// returns from the pure preflight before Deploy touches the cluster, so no fake +// clientset is needed. func TestDeploy_KafkaSASLPreflight(t *testing.T) { kafkaTrigger := &fn.ScaleOptions{ KEDA: &fn.KEDAScaleOptions{Triggers: []fn.KEDATrigger{{Type: "kafka"}}}, @@ -369,15 +369,6 @@ func TestDeploy_KafkaSASLPreflight(t *testing.T) { kafka *fn.KafkaConfig wantErr string }{ - { - name: "SASL_SSL with empty mechanism", - kafka: &fn.KafkaConfig{ - Brokers: "b:9092", Topic: "t", ConsumerGroup: "g", - SecurityProtocol: "SASL_SSL", - SASL: &fn.KafkaSASL{User: "u", Password: "p"}, - }, - wantErr: "run.kafka.sasl.mechanism is required", - }, { name: "SASL block with non-SASL protocol", kafka: &fn.KafkaConfig{ diff --git a/pkg/keda/kafka_scaling.go b/pkg/keda/kafka_scaling.go index 28a088ea4f..1e705dd34e 100644 --- a/pkg/keda/kafka_scaling.go +++ b/pkg/keda/kafka_scaling.go @@ -453,8 +453,23 @@ func buildScaledObject(f fn.Function, trigger fn.KEDATrigger, deployment *v1.Dep } } - if kafka.SASL != nil && kafka.SASL.Mechanism != "" { - triggerMeta["sasl"] = kedaSASLType(kafka.SASL.Mechanism) + if kafka.SecurityProtocol == "SASL_PLAINTEXT" || kafka.SecurityProtocol == "SASL_SSL" { + // Gate sasl on securityProtocol, mirroring func-go's Kafka runtime, + // which enables SASL from KAFKA_SECURITY_PROTOCOL rather than the + // presence of a sasl block -- and matching the tls gate above. The + // scaler reads the same EffectiveMechanism the function's container is + // wired with, so the two authenticate identically; gating on a non-empty + // mechanism instead would leave the scaler connecting without SASL while + // the function does, failing its lag reads. Validation ties a SASL_* + // protocol to a sasl block, but guard the deref for a direct Deploy + // caller that bypasses it. kedaSASLType returns "" only for a mechanism + // validation already rejects; skip the key in that defensive case rather + // than emitting an empty sasl value. + if kafka.SASL != nil { + if saslType := kedaSASLType(kafka.SASL.EffectiveMechanism()); saslType != "" { + triggerMeta["sasl"] = saslType + } + } } triggerSpec := map[string]interface{}{ diff --git a/pkg/keda/kafka_scaling_test.go b/pkg/keda/kafka_scaling_test.go index 8077e8f533..bc313aaea8 100644 --- a/pkg/keda/kafka_scaling_test.go +++ b/pkg/keda/kafka_scaling_test.go @@ -881,6 +881,40 @@ func TestBuildScaledObject(t *testing.T) { } } +// TestBuildScaledObject_EmptyMechanismEmitsPlaintext covers an omitted SASL +// mechanism: EffectiveMechanism resolves it to PLAIN, so the ScaledObject must +// still emit sasl: plaintext. Gating the sasl metadata on a non-empty mechanism +// would leave KEDA connecting without SASL while the function authenticates, so +// its lag reads fail and it never scales. +func TestBuildScaledObject_EmptyMechanismEmitsPlaintext(t *testing.T) { + f := fn.Function{ + Name: "test-func", + Run: fn.RunSpec{ + Kafka: &fn.KafkaConfig{ + Brokers: "broker:9093", + Topic: "t", + ConsumerGroup: "g", + SecurityProtocol: "SASL_SSL", + // Mechanism deliberately omitted: a SASL/PLAIN broker config + // that relies on func-go's empty-mechanism default. + SASL: &fn.KafkaSASL{User: "admin", Password: "{{ secret:s:k }}"}, + }, + }, + } + trigger := fn.KEDATrigger{Type: "kafka"} + + so := buildScaledObject(f, trigger, testDeployment(), "default", 0, 10) + if so == nil { + t.Fatal("expected ScaledObject, got nil") + } + spec := so.Object["spec"].(map[string]interface{}) + trigger0 := spec["triggers"].([]interface{})[0].(map[string]interface{}) + meta := trigger0["metadata"].(map[string]interface{}) + if meta["sasl"] != "plaintext" { + t.Errorf("sasl = %v, want plaintext for an empty SASL mechanism", meta["sasl"]) + } +} + func TestBuildScaledObject_TLSFromSecurityProtocol(t *testing.T) { // SecurityProtocol: SSL with no explicit run.kafka.tls block (relying on // the system's CA trust store) must still enable KEDA's tls handshake -- @@ -1027,7 +1061,10 @@ func TestKedaSASLType(t *testing.T) { "SCRAM-SHA-256": "scram_sha256", "SCRAM-SHA-512": "scram_sha512", "PLAIN": "plaintext", - "UNKNOWN": "", + // Callers pass KafkaSASL.EffectiveMechanism(), which never yields "", + // so an empty mechanism is an unmapped input like any other. + "": "", + "UNKNOWN": "", } for in, want := range tests { if got := kedaSASLType(in); got != want { diff --git a/schema/func_yaml-schema.json b/schema/func_yaml-schema.json index 9f8f60072c..3a70a2878d 100644 --- a/schema/func_yaml-schema.json +++ b/schema/func_yaml-schema.json @@ -419,7 +419,7 @@ "SCRAM-SHA-512" ], "type": "string", - "description": "SASL mechanism: PLAIN SCRAM-SHA-256 or SCRAM-SHA-512" + "description": "SASL mechanism: PLAIN SCRAM-SHA-256 or SCRAM-SHA-512. Optional; defaults to PLAIN when unset." }, "user": { "type": "string",