Skip to content

Commit

Permalink
[Kafka] Adding awsOptsFn (#791)
Browse files Browse the repository at this point in the history
  • Loading branch information
Tang8330 authored Jul 12, 2024
1 parent 60851fe commit 57979a9
Showing 1 changed file with 4 additions and 4 deletions.
8 changes: 4 additions & 4 deletions lib/kafkalib/connection.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,7 @@ func (c Connection) Mechanism() Mechanism {
return Plain
}

func (c Connection) Dialer(ctx context.Context) (*kafka.Dialer, error) {
func (c Connection) Dialer(ctx context.Context, awsOptFns ...func(options *awsCfg.LoadOptions) error) (*kafka.Dialer, error) {
dialer := &kafka.Dialer{
Timeout: 10 * time.Second,
DualStack: true,
Expand All @@ -66,7 +66,7 @@ func (c Connection) Dialer(ctx context.Context) (*kafka.Dialer, error) {
dialer.TLS = &tls.Config{}
}
case AwsMskIam:
_awsCfg, err := awsCfg.LoadDefaultConfig(ctx)
_awsCfg, err := awsCfg.LoadDefaultConfig(ctx, awsOptFns...)
if err != nil {
return nil, fmt.Errorf("failed to load aws configuration: %w", err)
}
Expand All @@ -83,7 +83,7 @@ func (c Connection) Dialer(ctx context.Context) (*kafka.Dialer, error) {
return dialer, nil
}

func (c Connection) Transport(ctx context.Context) (*kafka.Transport, error) {
func (c Connection) Transport(ctx context.Context, awsOptFns ...func(options *awsCfg.LoadOptions) error) (*kafka.Transport, error) {
transport := &kafka.Transport{
DialTimeout: 10 * time.Second,
}
Expand All @@ -100,7 +100,7 @@ func (c Connection) Transport(ctx context.Context) (*kafka.Transport, error) {
transport.TLS = &tls.Config{}
}
case AwsMskIam:
_awsCfg, err := awsCfg.LoadDefaultConfig(ctx)
_awsCfg, err := awsCfg.LoadDefaultConfig(ctx, awsOptFns...)
if err != nil {
return nil, fmt.Errorf("failed to load AWS configuration: %w", err)
}
Expand Down

0 comments on commit 57979a9

Please sign in to comment.