diff --git a/internal/sink/kafka/backend.go b/internal/sink/kafka/backend.go index bab55fb0..76810356 100644 --- a/internal/sink/kafka/backend.go +++ b/internal/sink/kafka/backend.go @@ -6,6 +6,7 @@ package kafka import ( "context" "encoding/json" + "errors" "fmt" "strings" "time" @@ -143,7 +144,7 @@ func dialTransport(cfg Config) (*kafka.Transport, error) { if cfg.Username != "" { mechanism, err := scram.Mechanism(scram.SHA256, cfg.Username, cfg.Password) if err != nil { - return nil, fmt.Errorf("kafka SASL: %w", err) + return nil, redactedSASLError() } transport.SASL = mechanism @@ -151,3 +152,7 @@ func dialTransport(cfg Config) (*kafka.Transport, error) { return transport, nil } + +func redactedSASLError() error { + return errors.New("kafka SASL: authentication configuration failed") +} diff --git a/internal/sink/kafka/sasl_redact_test.go b/internal/sink/kafka/sasl_redact_test.go new file mode 100644 index 00000000..4ab2d5cb --- /dev/null +++ b/internal/sink/kafka/sasl_redact_test.go @@ -0,0 +1,27 @@ +// SPDX-License-Identifier: MIT +// Copyright (c) 2026 Konrad Heimel + +package kafka + +import ( + "strings" + "testing" +) + +func TestDialTransport_SASLPrepErrorOmitsPassword(t *testing.T) { + t.Parallel() + + const secret = "s3cret\x07" + _, err := dialTransport(Config{ + Brokers: []string{"127.0.0.1:1"}, + Topic: "inventory", + Username: "alice", + Password: secret, + }) + if err == nil { + t.Fatal("expected SASL construction error for a prohibited password rune") + } + if strings.Contains(err.Error(), "s3cret") || strings.Contains(err.Error(), secret) { + t.Fatalf("SASL error must not include the password: %q", err.Error()) + } +}