Skip to content
Closed
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
21 changes: 1 addition & 20 deletions aws/keyspaces/keyspaces.go
Original file line number Diff line number Diff line change
Expand Up @@ -215,26 +215,7 @@ func getConfig() (aws.Config, error) {
return aws.Config{}, errors.New("AWS_REGION is not set")
}

var cfg aws.Config
var err error

if awsEndpoint := os.Getenv("AWS_ENDPOINT_URL"); awsEndpoint != "" {
customResolver := aws.EndpointResolverWithOptionsFunc(func(service, region string, options ...interface{}) (aws.Endpoint, error) {
return aws.Endpoint{
PartitionID: "aws",
URL: awsEndpoint,
HostnameImmutable: true,
}, nil
})

cfg, err = config.LoadDefaultConfig(
context.TODO(),
config.WithEndpointResolverWithOptions(customResolver),
)
} else {
cfg, err = config.LoadDefaultConfig(context.TODO())
}

cfg, err := config.LoadDefaultConfig(context.TODO())
if err != nil {
return aws.Config{}, err
}
Expand Down
59 changes: 38 additions & 21 deletions aws/s3/s3.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ import (
"github.com/aws/aws-sdk-go-v2/service/s3"
"github.com/aws/aws-sdk-go-v2/service/s3/types"
"github.com/aws/smithy-go"
smithyendpoints "github.com/aws/smithy-go/endpoints"
)

type S3 struct {
Expand All @@ -36,7 +37,7 @@ func New() (S3, error) {
if err != nil {
return S3{}, err
}
return S3{client: s3.NewFromConfig(cfg)}, nil
return S3{client: newFromConfig(cfg)}, nil
}

// NewWithMaxRetries returns the same as New(), but with the
Expand All @@ -46,7 +47,7 @@ func NewWithMaxRetries(maxRetries int) (S3, error) {
if err != nil {
return S3{}, err
}
client := s3.NewFromConfig(cfg, func(options *s3.Options) {
client := newFromConfig(cfg, func(options *s3.Options) {
options.Retryer = retry.AddWithMaxAttempts(options.Retryer, maxRetries)
})
return S3{client: client}, nil
Expand All @@ -59,7 +60,7 @@ func NewWithOptions(optFns ...func(*s3.Options)) (S3, error) {
if err != nil {
return S3{}, err
}
client := s3.NewFromConfig(cfg, optFns...)
client := newFromConfig(cfg, optFns...)
return S3{client: client}, nil
}

Expand Down Expand Up @@ -90,29 +91,45 @@ func getConfig() (aws.Config, error) {
return aws.Config{}, errors.New("AWS_REGION is not set")
}

var cfg aws.Config
var err error
cfg, err := config.LoadDefaultConfig(context.TODO())
if err != nil {
return aws.Config{}, err
}
return cfg, nil
}

if awsEndpoint := os.Getenv("AWS_ENDPOINT_URL"); awsEndpoint != "" {
customResolver := aws.EndpointResolverWithOptionsFunc(func(service, region string, options ...interface{}) (aws.Endpoint, error) {
return aws.Endpoint{
PartitionID: "aws",
URL: awsEndpoint,
HostnameImmutable: true,
}, nil
})
// newFromConfig adds the custom endpoint option to the list of options given (if
// the AWS_ENDPOINT_URL env var is set), and returns a Client.
func newFromConfig(config aws.Config, optFns ...func(*s3.Options)) *s3.Client {

cfg, err = config.LoadDefaultConfig(
context.TODO(),
config.WithEndpointResolverWithOptions(customResolver))
} else {
cfg, err = config.LoadDefaultConfig(context.TODO())
if awsEndpoint := os.Getenv("AWS_ENDPOINT_URL"); awsEndpoint != "" {
customResolver := func(options *s3.Options) {
options.EndpointResolverV2 = newCustomEndpointResolver(awsEndpoint)
options.UsePathStyle = true
}
optFns = append(optFns, customResolver)
}

if err != nil {
return aws.Config{}, err
client := s3.NewFromConfig(config, optFns...)

return client
}

type customEndpointResolver struct {
inner s3.EndpointResolverV2
endpoint *string
}

func newCustomEndpointResolver(endpoint string) *customEndpointResolver {
return &customEndpointResolver{
inner: s3.NewDefaultEndpointResolverV2(),
endpoint: &endpoint,
}
return cfg, nil
}

func (c *customEndpointResolver) ResolveEndpoint(ctx context.Context, params s3.EndpointParameters) (smithyendpoints.Endpoint, error) {
params.Endpoint = c.endpoint
return c.inner.ResolveEndpoint(ctx, params)
}

// Ready returns whether the S3 client has been initialised.
Expand Down
56 changes: 36 additions & 20 deletions aws/sns/sns.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"github.com/aws/aws-sdk-go-v2/aws/retry"
"github.com/aws/aws-sdk-go-v2/config"
"github.com/aws/aws-sdk-go-v2/service/sns"
smithyendpoints "github.com/aws/smithy-go/endpoints"
)

type SNS struct {
Expand All @@ -25,7 +26,7 @@ func New() (SNS, error) {
if err != nil {
return SNS{}, err
}
return SNS{client: sns.NewFromConfig(cfg)}, nil
return SNS{client: newFromConfig(cfg)}, nil
}

// NewWithMaxRetries returns the same as New(), but with the
Expand All @@ -35,7 +36,7 @@ func NewWithMaxRetries(maxRetries int) (SNS, error) {
if err != nil {
return SNS{}, err
}
client := sns.NewFromConfig(cfg, func(options *sns.Options) {
client := newFromConfig(cfg, func(options *sns.Options) {
options.Retryer = retry.AddWithMaxAttempts(options.Retryer, maxRetries)
})
return SNS{client: client}, nil
Expand All @@ -47,29 +48,44 @@ func getConfig() (aws.Config, error) {
return aws.Config{}, errors.New("AWS_REGION is not set")
}

var cfg aws.Config
var err error
cfg, err := config.LoadDefaultConfig(context.TODO())
if err != nil {
return aws.Config{}, err
}
return cfg, nil
}

// newFromConfig adds the custom endpoint option to the list of options given (if
// the AWS_ENDPOINT_URL env var is set), and returns a Client.
func newFromConfig(config aws.Config, optFns ...func(*sns.Options)) *sns.Client {

if awsEndpoint := os.Getenv("AWS_ENDPOINT_URL"); awsEndpoint != "" {
customResolver := aws.EndpointResolverWithOptionsFunc(func(service, region string, options ...interface{}) (aws.Endpoint, error) {
return aws.Endpoint{
PartitionID: "aws",
SigningRegion: region,
URL: awsEndpoint,
}, nil
})

cfg, err = config.LoadDefaultConfig(
context.TODO(),
config.WithEndpointResolverWithOptions(customResolver))
} else {
cfg, err = config.LoadDefaultConfig(context.TODO())
customResolver := func(options *sns.Options) {
options.EndpointResolverV2 = newCustomEndpointResolver(awsEndpoint)
}
optFns = append(optFns, customResolver)
}

if err != nil {
return aws.Config{}, err
client := sns.NewFromConfig(config, optFns...)

return client
}

type customEndpointResolver struct {
inner sns.EndpointResolverV2
endpoint *string
}

func newCustomEndpointResolver(endpoint string) *customEndpointResolver {
return &customEndpointResolver{
inner: sns.NewDefaultEndpointResolverV2(),
endpoint: &endpoint,
}
return cfg, nil
}

func (c *customEndpointResolver) ResolveEndpoint(ctx context.Context, params sns.EndpointParameters) (smithyendpoints.Endpoint, error) {
params.Endpoint = c.endpoint
return c.inner.ResolveEndpoint(ctx, params)
}

// Ready returns whether the SNS client has been initialised.
Expand Down
56 changes: 36 additions & 20 deletions aws/sqs/sqs.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (
"github.com/aws/aws-sdk-go-v2/service/sqs"
"github.com/aws/aws-sdk-go-v2/service/sqs/types"
smithy "github.com/aws/smithy-go"
smithyendpoints "github.com/aws/smithy-go/endpoints"
)

type Raw struct {
Expand All @@ -39,7 +40,7 @@ func New() (SQS, error) {
if err != nil {
return SQS{}, err
}
return SQS{client: sqs.NewFromConfig(cfg)}, nil
return SQS{client: newFromConfig(cfg)}, nil
}

// NewWithMaxRetries returns the same as New(), but with the
Expand All @@ -49,7 +50,7 @@ func NewWithMaxRetries(maxRetries int) (SQS, error) {
if err != nil {
return SQS{}, err
}
client := sqs.NewFromConfig(cfg, func(options *sqs.Options) {
client := newFromConfig(cfg, func(options *sqs.Options) {
options.Retryer = retry.AddWithMaxAttempts(options.Retryer, maxRetries)
})

Expand All @@ -62,29 +63,44 @@ func getConfig() (aws.Config, error) {
return aws.Config{}, errors.New("AWS_REGION is not set")
}

var cfg aws.Config
var err error
cfg, err := config.LoadDefaultConfig(context.TODO())
if err != nil {
return aws.Config{}, err
}
return cfg, nil
}

if awsEndpoint := os.Getenv("AWS_ENDPOINT_URL"); awsEndpoint != "" {
customResolver := aws.EndpointResolverWithOptionsFunc(func(service, region string, options ...interface{}) (aws.Endpoint, error) {
return aws.Endpoint{
PartitionID: "aws",
SigningRegion: region,
URL: awsEndpoint,
}, nil
})
// newFromConfig adds the custom endpoint option to the list of options given (if
// the AWS_ENDPOINT_URL env var is set), and returns a Client.
func newFromConfig(config aws.Config, optFns ...func(*sqs.Options)) *sqs.Client {

cfg, err = config.LoadDefaultConfig(
context.TODO(),
config.WithEndpointResolverWithOptions(customResolver))
} else {
cfg, err = config.LoadDefaultConfig(context.TODO())
if awsEndpoint := os.Getenv("AWS_ENDPOINT_URL"); awsEndpoint != "" {
customResolver := func(options *sqs.Options) {
options.EndpointResolverV2 = newCustomEndpointResolver(awsEndpoint)
}
optFns = append(optFns, customResolver)
}

if err != nil {
return aws.Config{}, err
client := sqs.NewFromConfig(config, optFns...)

return client
}

type customEndpointResolver struct {
inner sqs.EndpointResolverV2
endpoint *string
}

func newCustomEndpointResolver(endpoint string) *customEndpointResolver {
return &customEndpointResolver{
inner: sqs.NewDefaultEndpointResolverV2(),
endpoint: &endpoint,
}
return cfg, nil
}

func (c *customEndpointResolver) ResolveEndpoint(ctx context.Context, params sqs.EndpointParameters) (smithyendpoints.Endpoint, error) {
params.Endpoint = c.endpoint
return c.inner.ResolveEndpoint(ctx, params)
}

// Ready returns whether the SQS client has been initialised.
Expand Down
26 changes: 13 additions & 13 deletions go.mod
Original file line number Diff line number Diff line change
@@ -1,37 +1,37 @@
module github.com/GeoNet/kit

go 1.23
go 1.25.0

require (
github.com/aws/aws-sdk-go-v2 v1.25.3
github.com/aws/aws-sdk-go-v2 v1.43.3
github.com/aws/aws-sdk-go-v2/config v1.27.7
github.com/aws/aws-sdk-go-v2/feature/s3/manager v1.16.11
github.com/aws/aws-sdk-go-v2/service/s3 v1.52.1
github.com/aws/aws-sdk-go-v2/service/s3 v1.106.4
github.com/aws/aws-sdk-go-v2/service/sns v1.29.2
github.com/aws/aws-sdk-go-v2/service/sqs v1.31.2
github.com/aws/aws-sigv4-auth-cassandra-gocql-driver-plugin v1.1.0
github.com/aws/smithy-go v1.20.1
github.com/aws/smithy-go v1.27.6
github.com/gocql/gocql v1.6.0
github.com/golang/groupcache v0.0.0-20210331224755-41bb18bfe9da
github.com/lib/pq v1.3.0
github.com/stretchr/testify v1.11.1
golang.org/x/text v0.14.0
golang.org/x/text v0.40.0
google.golang.org/protobuf v1.33.0
)

require (
github.com/aws/aws-sdk-go v1.51.1 // indirect
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.6.1 // indirect
github.com/aws/aws-sdk-go-v2/aws/protocol/eventstream v1.7.16 // indirect
github.com/aws/aws-sdk-go-v2/credentials v1.17.7 // indirect
github.com/aws/aws-sdk-go-v2/feature/ec2/imds v1.15.3 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.3.3 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.6.3 // indirect
github.com/aws/aws-sdk-go-v2/internal/configsources v1.4.34 // indirect
github.com/aws/aws-sdk-go-v2/internal/endpoints/v2 v2.7.34 // indirect
github.com/aws/aws-sdk-go-v2/internal/ini v1.8.0 // indirect
github.com/aws/aws-sdk-go-v2/internal/v4a v1.3.3 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.11.1 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.3.5 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.11.5 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.17.3 // indirect
github.com/aws/aws-sdk-go-v2/internal/v4a v1.4.35 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/accept-encoding v1.13.15 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.27 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.34 // indirect
github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.35 // indirect
github.com/aws/aws-sdk-go-v2/service/sso v1.20.2 // indirect
github.com/aws/aws-sdk-go-v2/service/ssooidc v1.23.2 // indirect
github.com/aws/aws-sdk-go-v2/service/sts v1.28.4 // indirect
Expand Down
Loading