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
2 changes: 1 addition & 1 deletion docs/reference/func_yaml.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).

Expand Down
11 changes: 6 additions & 5 deletions pkg/functions/function.go
Original file line number Diff line number Diff line change
Expand Up @@ -276,11 +276,12 @@ 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: func-go's Kafka runtime defaults an empty
// mechanism to PLAIN and kedaSASLType maps "" to KEDA's "plaintext", so
// the function container and the scaler 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")
Expand Down
11 changes: 7 additions & 4 deletions pkg/functions/function_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -743,17 +743,20 @@ 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: func-go's runtime defaults it to
// PLAIN and the keda scaler maps it to "plaintext", 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",
ConsumerGroup: "my-group",
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",
Expand Down
9 changes: 0 additions & 9 deletions pkg/keda/deployer_unit_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{
Expand Down
20 changes: 17 additions & 3 deletions pkg/keda/kafka_scaling.go
Original file line number Diff line number Diff line change
Expand Up @@ -405,7 +405,10 @@ func kedaSASLType(mechanism string) string {
return "scram_sha256"
case "SCRAM-SHA-512":
return "scram_sha512"
case "PLAIN":
case "PLAIN", "":
// An empty mechanism defaults to PLAIN in func-go's Kafka runtime, so
// the scaler must authenticate the same way for its lag reads to match
// the function's consumption.
Comment thread
aliok marked this conversation as resolved.
return "plaintext"
default:
return ""
Expand Down Expand Up @@ -453,8 +456,19 @@ 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.SASL != nil {
// Emit sasl whenever the SASL block is configured, not only when a
// mechanism is named: func-go's Kafka runtime treats a present SASL
// block as SASL-enabled and defaults an empty mechanism to PLAIN, so
// kedaSASLType maps "" to "plaintext". Gating on Mechanism != "" here
// would leave the scaler connecting without SASL while the function
// authenticates, so its lag reads fail and the ScaledObject never
// scales. kedaSASLType returns "" only for a mechanism that validation
// already rejects; skip the key in that defensive case rather than
// emitting an empty sasl value.
if saslType := kedaSASLType(kafka.SASL.Mechanism); saslType != "" {
triggerMeta["sasl"] = saslType
}
}

triggerSpec := map[string]interface{}{
Expand Down
39 changes: 38 additions & 1 deletion pkg/keda/kafka_scaling_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -881,6 +881,40 @@ func TestBuildScaledObject(t *testing.T) {
}
}

// TestBuildScaledObject_EmptyMechanismEmitsPlaintext covers the empty-mechanism
// SASL config this PR newly accepts: func-go defaults an omitted mechanism to
// SASL/PLAIN, so the ScaledObject must still emit sasl: plaintext. Gating the
// sasl metadata on 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 --
Expand Down Expand Up @@ -1027,7 +1061,10 @@ func TestKedaSASLType(t *testing.T) {
"SCRAM-SHA-256": "scram_sha256",
"SCRAM-SHA-512": "scram_sha512",
"PLAIN": "plaintext",
"UNKNOWN": "",
// An empty mechanism defaults to PLAIN in func-go, so the scaler must
// use the matching "plaintext" type rather than an empty value.
"": "plaintext",
"UNKNOWN": "",
}
for in, want := range tests {
if got := kedaSASLType(in); got != want {
Expand Down
Loading