Skip to content

Commit a163fbd

Browse files
committed
Add support for role chain assumption
1 parent bd8aade commit a163fbd

2 files changed

Lines changed: 16 additions & 2 deletions

File tree

internal/clients/kafka/client.go

Lines changed: 15 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,7 +11,10 @@ import (
1111
"strings"
1212
"time"
1313

14+
"github.com/aws/aws-sdk-go-v2/aws"
1415
"github.com/aws/aws-sdk-go-v2/config"
16+
"github.com/aws/aws-sdk-go-v2/credentials/stscreds"
17+
"github.com/aws/aws-sdk-go-v2/service/sts"
1518
"github.com/twmb/franz-go/pkg/kadm"
1619
"github.com/twmb/franz-go/pkg/kgo"
1720
"github.com/twmb/franz-go/pkg/sasl"
@@ -70,7 +73,9 @@ func NewAdminClient(ctx context.Context, data []byte, kube client.Client) (*kadm
7073
Pass: kc.SASL.Password,
7174
}.AsMechanism()
7275
case "aws-msk-iam":
73-
mechanism = kaws.ManagedStreamingIAM(authenticateAwsIam)
76+
mechanism = kaws.ManagedStreamingIAM(func(ctx context.Context) (kaws.Auth, error) {
77+
return authenticateAwsIam(ctx, kc.SASL.RoleArn)
78+
})
7479
opts = append(opts, kgo.Dialer((&tls.Dialer{NetDialer: &net.Dialer{Timeout: 10 * time.Second}}).DialContext))
7580
case "scram-sha-512":
7681
mechanism = scram.Auth{
@@ -99,12 +104,20 @@ func NewAdminClient(ctx context.Context, data []byte, kube client.Client) (*kadm
99104
return kadm.NewClient(c), nil
100105
}
101106

102-
func authenticateAwsIam(ctx context.Context) (a kaws.Auth, err error) {
107+
func authenticateAwsIam(ctx context.Context, roleArn string) (a kaws.Auth, err error) {
103108
s, err := config.LoadDefaultConfig(ctx)
104109
if err != nil {
105110
return kaws.Auth{}, err
106111
}
107112

113+
if roleArn != "" {
114+
stsClient := sts.NewFromConfig(s)
115+
provider := stscreds.NewAssumeRoleProvider(stsClient, roleArn, func(o *stscreds.AssumeRoleOptions) {
116+
o.RoleSessionName = "crossplane-provider-kafka"
117+
})
118+
s.Credentials = aws.NewCredentialsCache(provider)
119+
}
120+
108121
v, err := s.Credentials.Retrieve(ctx)
109122
if err != nil {
110123
return kaws.Auth{}, err

internal/clients/kafka/config.go

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -10,6 +10,7 @@ type Config struct {
1010
// SASL is an sasl option
1111
type SASL struct {
1212
Mechanism string `json:"mechanism"`
13+
RoleArn string `json:"roleArn"`
1314
Username string `json:"username"`
1415
Password string `json:"password"` //nolint:gosec
1516
}

0 commit comments

Comments
 (0)