Skip to content

Commit a79cc8f

Browse files
committed
Migrate EC2 metadata to SDKv2 (#1992)
1 parent 7d84436 commit a79cc8f

24 files changed

Lines changed: 730 additions & 313 deletions

extension/agenthealth/handler/stats/agent/flag.go

Lines changed: 7 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -108,6 +108,8 @@ type FlagSet interface {
108108
SetValues(flags map[Flag]any)
109109
// OnChange registers a callback that triggers on flag sets.
110110
OnChange(callback func())
111+
// Reset allows tests to clear the stored flags
112+
Reset()
111113
}
112114

113115
type flagSet struct {
@@ -180,6 +182,11 @@ func (p *flagSet) notify() {
180182
}
181183
}
182184

185+
func (p *flagSet) Reset() {
186+
p.m.Clear()
187+
p.notify()
188+
}
189+
183190
func UsageFlags() FlagSet {
184191
flagOnce.Do(func() {
185192
flagSingleton = &flagSet{}

extension/entitystore/ec2Info_test.go

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -9,22 +9,22 @@ import (
99
"testing"
1010
"time"
1111

12-
"github.com/aws/aws-sdk-go/aws/ec2metadata"
12+
"github.com/aws/aws-sdk-go-v2/feature/ec2/imds"
1313
"github.com/stretchr/testify/assert"
1414
"go.uber.org/zap"
1515

1616
"github.com/aws/amazon-cloudwatch-agent/internal/ec2metadataprovider"
1717
)
1818

19-
var mockedInstanceIdentityDoc = &ec2metadata.EC2InstanceIdentityDocument{
19+
var mockedInstanceIdentityDoc = &imds.InstanceIdentityDocument{
2020
InstanceID: "i-01d2417c27a396e44",
2121
AccountID: "874389809020",
2222
Region: "us-east-1",
2323
InstanceType: "m5ad.large",
2424
ImageID: "ami-09edd32d9b0990d49",
2525
}
2626

27-
var mockedInstanceIdentityDocWithLargeInstanceId = &ec2metadata.EC2InstanceIdentityDocument{
27+
var mockedInstanceIdentityDocWithLargeInstanceID = &imds.InstanceIdentityDocument{
2828
InstanceID: "i-01d2417c27a396e44394824728",
2929
AccountID: "874389809020",
3030
Region: "us-east-1",
@@ -60,12 +60,12 @@ func TestSetInstanceIDAccountID(t *testing.T) {
6060
{
6161
name: "InstanceId too large",
6262
args: args{
63-
metadataProvider: &mockMetadataProvider{InstanceIdentityDocument: mockedInstanceIdentityDocWithLargeInstanceId},
63+
metadataProvider: &mockMetadataProvider{InstanceIdentityDocument: mockedInstanceIdentityDocWithLargeInstanceID},
6464
},
6565
wantErr: false,
6666
want: EC2Info{
6767
InstanceID: "",
68-
AccountID: mockedInstanceIdentityDocWithLargeInstanceId.AccountID,
68+
AccountID: mockedInstanceIdentityDocWithLargeInstanceID.AccountID,
6969
},
7070
},
7171
}

extension/entitystore/extension.go

Lines changed: 13 additions & 40 deletions
Original file line numberDiff line numberDiff line change
@@ -7,20 +7,17 @@ import (
77
"context"
88
"time"
99

10+
"github.com/aws/aws-sdk-go-v2/aws"
1011
"github.com/aws/aws-sdk-go-v2/service/cloudwatchlogs/types"
11-
"github.com/aws/aws-sdk-go/aws"
12-
"github.com/aws/aws-sdk-go/aws/client"
13-
"github.com/aws/aws-sdk-go/service/ec2"
14-
"github.com/aws/aws-sdk-go/service/ec2/ec2iface"
1512
"github.com/jellydator/ttlcache/v3"
1613
"go.opentelemetry.io/collector/component"
1714
"go.opentelemetry.io/collector/extension"
1815
"go.uber.org/atomic"
1916
"go.uber.org/zap"
2017

21-
configaws "github.com/aws/amazon-cloudwatch-agent/cfg/aws"
18+
configaws "github.com/aws/amazon-cloudwatch-agent/cfg/aws/v2"
2219
"github.com/aws/amazon-cloudwatch-agent/internal/ec2metadataprovider"
23-
"github.com/aws/amazon-cloudwatch-agent/internal/retryer"
20+
"github.com/aws/amazon-cloudwatch-agent/internal/retryer/v2"
2421
"github.com/aws/amazon-cloudwatch-agent/plugins/processors/awsentity/entityattributes"
2522
"github.com/aws/amazon-cloudwatch-agent/translator/config"
2623
)
@@ -35,8 +32,6 @@ const (
3532
podTerminationCheckInterval = 5 * time.Minute
3633
)
3734

38-
type ec2ProviderType func(string, *configaws.CredentialConfig) ec2iface.EC2API
39-
4035
type serviceProviderInterface interface {
4136
startServiceProvider()
4237
addEntryForLogFile(LogFileGlob, ServiceAttribute)
@@ -69,31 +64,23 @@ type EntityStore struct {
6964
// that we can attach to the entity
7065
serviceprovider serviceProviderInterface
7166

72-
// nativeCredential stores the credential config for agent's native
73-
// component such as LogAgent
74-
nativeCredential client.ConfigProvider
75-
7667
metadataprovider ec2metadataprovider.MetadataProvider
7768

7869
podTerminationCheckInterval time.Duration
7970
}
8071

8172
var _ extension.Extension = (*EntityStore)(nil)
8273

83-
func (e *EntityStore) Start(ctx context.Context, host component.Host) error {
74+
func (e *EntityStore) Start(ctx context.Context, _ component.Host) error {
8475
// Get IMDS client and EC2 API client which requires region for authentication
8576
// These will be passed down to any object that requires access to IMDS or EC2
8677
// API client so we have single source of truth for credential
8778
e.done = make(chan struct{})
88-
e.metadataprovider = getMetaDataProvider()
79+
e.metadataprovider = getMetaDataProvider(ctx)
8980
e.mode = e.config.Mode
9081
e.kubernetesMode = e.config.KubernetesMode
9182
e.podTerminationCheckInterval = podTerminationCheckInterval
92-
ec2CredentialConfig := &configaws.CredentialConfig{
93-
Profile: e.config.Profile,
94-
Filename: e.config.Filename,
95-
}
96-
e.serviceprovider = newServiceProvider(e.mode, e.config.Region, &e.ec2Info, e.metadataprovider, getEC2Provider, ec2CredentialConfig, e.done, e.logger)
83+
e.serviceprovider = newServiceProvider(e.mode, e.config.Region, &e.ec2Info, e.metadataprovider, e.done, e.logger)
9784
switch e.mode {
9885
case config.ModeEC2:
9986
e.ec2Info = *newEC2Info(e.metadataprovider, e.done, e.config.Region, e.logger)
@@ -138,14 +125,6 @@ func (e *EntityStore) EC2Info() EC2Info {
138125
return e.ec2Info
139126
}
140127

141-
func (e *EntityStore) SetNativeCredential(client client.ConfigProvider) {
142-
e.nativeCredential = client
143-
}
144-
145-
func (e *EntityStore) NativeCredentialExists() bool {
146-
return e.nativeCredential != nil
147-
}
148-
149128
// CreateLogFileEntity creates the entity for log events that are being uploaded from a log file in the environment.
150129
func (e *EntityStore) CreateLogFileEntity(logFileGlob LogFileGlob, logGroupName LogGroupName) *types.Entity {
151130
if e.serviceprovider == nil {
@@ -260,19 +239,13 @@ func (e *EntityStore) createServiceKeyAttributes(serviceAttr ServiceAttribute) m
260239
return serviceKeyAttr
261240
}
262241

263-
var getMetaDataProvider = func() ec2metadataprovider.MetadataProvider {
264-
mdCredentialConfig := &configaws.CredentialConfig{}
265-
return ec2metadataprovider.NewMetadataProvider(mdCredentialConfig.Credentials(), retryer.GetDefaultRetryNumber())
266-
}
267-
268-
var getEC2Provider = func(region string, ec2CredentialConfig *configaws.CredentialConfig) ec2iface.EC2API {
269-
ec2CredentialConfig.Region = region
270-
return ec2.New(
271-
ec2CredentialConfig.Credentials(),
272-
&aws.Config{
273-
LogLevel: configaws.SDKLogLevel(),
274-
Logger: configaws.SDKLogger{},
275-
})
242+
var getMetaDataProvider = func(ctx context.Context) ec2metadataprovider.MetadataProvider {
243+
mdCredentialConfig := &configaws.CredentialsConfig{}
244+
cfg, err := mdCredentialConfig.LoadConfig(ctx)
245+
if err != nil {
246+
cfg = aws.Config{}
247+
}
248+
return ec2metadataprovider.NewMetadataProvider(cfg, retryer.GetDefaultRetryNumber())
276249
}
277250

278251
func addNonEmptyToMap(m map[string]string, key, value string) {

extension/entitystore/extension_test.go

Lines changed: 19 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -12,9 +12,8 @@ import (
1212
"testing"
1313
"time"
1414

15+
"github.com/aws/aws-sdk-go-v2/feature/ec2/imds"
1516
"github.com/aws/aws-sdk-go-v2/service/cloudwatchlogs/types"
16-
"github.com/aws/aws-sdk-go/aws/ec2metadata"
17-
"github.com/aws/aws-sdk-go/aws/session"
1817
"github.com/jellydator/ttlcache/v3"
1918
"github.com/stretchr/testify/assert"
2019
"github.com/stretchr/testify/mock"
@@ -73,57 +72,59 @@ func (s *mockServiceProvider) setAutoScalingGroup(asg string) {
7372
}
7473

7574
type mockMetadataProvider struct {
76-
InstanceIdentityDocument *ec2metadata.EC2InstanceIdentityDocument
75+
InstanceIdentityDocument *imds.InstanceIdentityDocument
7776
Tags map[string]string
7877
InstanceTagError bool
7978
}
8079

81-
func mockMetadataProviderFunc() ec2metadataprovider.MetadataProvider {
80+
var _ ec2metadataprovider.MetadataProvider = (*mockMetadataProvider)(nil)
81+
82+
func mockMetadataProviderFunc(context.Context) ec2metadataprovider.MetadataProvider {
8283
return &mockMetadataProvider{
8384
Tags: map[string]string{
8485
"aws:autoscaling:groupName": "ASG-1",
8586
},
86-
InstanceIdentityDocument: &ec2metadata.EC2InstanceIdentityDocument{
87+
InstanceIdentityDocument: &imds.InstanceIdentityDocument{
8788
InstanceID: "i-123456789",
8889
},
8990
}
9091
}
9192

9293
func mockMetadataProviderWithAccountId(accountId string) *mockMetadataProvider {
9394
return &mockMetadataProvider{
94-
InstanceIdentityDocument: &ec2metadata.EC2InstanceIdentityDocument{
95+
InstanceIdentityDocument: &imds.InstanceIdentityDocument{
9596
AccountID: accountId,
9697
},
9798
}
9899
}
99100

100-
func (m *mockMetadataProvider) Get(ctx context.Context) (ec2metadata.EC2InstanceIdentityDocument, error) {
101+
func (m *mockMetadataProvider) Get(context.Context) (imds.InstanceIdentityDocument, error) {
101102
if m.InstanceIdentityDocument != nil {
102103
return *m.InstanceIdentityDocument, nil
103104
}
104-
return ec2metadata.EC2InstanceIdentityDocument{}, errors.New("No instance identity document")
105+
return imds.InstanceIdentityDocument{}, errors.New("no instance identity document")
105106
}
106107

107-
func (m *mockMetadataProvider) Hostname(ctx context.Context) (string, error) {
108+
func (m *mockMetadataProvider) Hostname(context.Context) (string, error) {
108109
return "MockHostName", nil
109110
}
110111

111-
func (m *mockMetadataProvider) InstanceID(ctx context.Context) (string, error) {
112+
func (m *mockMetadataProvider) InstanceID(context.Context) (string, error) {
112113
return "MockInstanceID", nil
113114
}
114115

115-
func (m *mockMetadataProvider) InstanceTags(_ context.Context) ([]string, error) {
116+
func (m *mockMetadataProvider) InstanceTags(context.Context) ([]string, error) {
116117
if m.InstanceTagError {
117118
return nil, errors.New("an error occurred for instance tag retrieval")
118119
}
119120
return maps.Keys(m.Tags), nil
120121
}
121122

122-
func (m *mockMetadataProvider) ClientIAMRole(ctx context.Context) (string, error) {
123+
func (m *mockMetadataProvider) ClientIAMRole(context.Context) (string, error) {
123124
return "TestRole", nil
124125
}
125126

126-
func (m *mockMetadataProvider) InstanceTagValue(ctx context.Context, tagKey string) (string, error) {
127+
func (m *mockMetadataProvider) InstanceTagValue(_ context.Context, tagKey string) (string, error) {
127128
tag, ok := m.Tags[tagKey]
128129
if !ok {
129130
return "", errors.New("tag not found")
@@ -333,10 +334,9 @@ func TestEntityStore_createLogFileRID(t *testing.T) {
333334
sp.On("logFileServiceAttribute", glob, group).Return(serviceAttr)
334335
sp.On("getAutoScalingGroup").Return("ASG-1")
335336
e := EntityStore{
336-
mode: config.ModeEC2,
337-
ec2Info: EC2Info{InstanceID: instanceId, AccountID: accountId},
338-
serviceprovider: sp,
339-
nativeCredential: &session.Session{},
337+
mode: config.ModeEC2,
338+
ec2Info: EC2Info{InstanceID: instanceId, AccountID: accountId},
339+
serviceprovider: sp,
340340
}
341341

342342
entity := e.CreateLogFileEntity(glob, group)
@@ -364,9 +364,8 @@ func TestEntityStore_createLogFileRID_ServiceProviderIsEmpty(t *testing.T) {
364364
glob := LogFileGlob("glob")
365365
group := LogGroupName("group")
366366
e := EntityStore{
367-
mode: config.ModeEC2,
368-
ec2Info: EC2Info{InstanceID: instanceId},
369-
nativeCredential: &session.Session{},
367+
mode: config.ModeEC2,
368+
ec2Info: EC2Info{InstanceID: instanceId},
370369
}
371370

372371
entity := e.CreateLogFileEntity(glob, group)
@@ -548,7 +547,6 @@ func TestEntityStore_GetMetricServiceNameSource(t *testing.T) {
548547
ec2Info: EC2Info{InstanceID: instanceId},
549548
serviceprovider: sp,
550549
metadataprovider: mockMetadataProviderWithAccountId(accountId),
551-
nativeCredential: &session.Session{},
552550
}
553551

554552
serviceName, serviceNameSource := e.GetMetricServiceNameAndSource()
@@ -564,7 +562,6 @@ func TestEntityStore_GetMetricServiceNameSource_ServiceProviderEmpty(t *testing.
564562
mode: config.ModeEC2,
565563
ec2Info: EC2Info{InstanceID: instanceId},
566564
metadataprovider: mockMetadataProviderWithAccountId(accountId),
567-
nativeCredential: &session.Session{},
568565
}
569566

570567
serviceName, serviceNameSource := e.GetMetricServiceNameAndSource()

extension/entitystore/retryer.go

Lines changed: 13 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -4,16 +4,16 @@
44
package entitystore
55

66
import (
7+
"errors"
78
"math/rand"
89
"time"
910

10-
"github.com/aws/aws-sdk-go/aws/awserr"
11+
"github.com/aws/smithy-go"
1112
"go.uber.org/zap"
1213
)
1314

1415
const (
15-
RequestLimitExceeded = "RequestLimitExceeded"
16-
infRetry = -1
16+
infRetry = -1
1717
)
1818

1919
var (
@@ -57,7 +57,7 @@ func (r *Retryer) refreshLoop(updateFunc func() error) int {
5757
err := updateFunc()
5858
if err == nil && r.oneTime {
5959
return retry
60-
} else if awsErr, ok := err.(awserr.Error); ok && !r.retryAnyError && !retryableErrorMap[awsErr.Code()] {
60+
} else if !r.retryAnyError && !isRetryableError(err) {
6161
return retry
6262
}
6363

@@ -83,7 +83,15 @@ func (r *Retryer) refreshLoop(updateFunc func() error) int {
8383
}
8484

8585
}
86-
return retry
86+
}
87+
88+
// isRetryableError checks if the error is a retryable API error code recognized by the extension.
89+
func isRetryableError(err error) bool {
90+
var apiErr smithy.APIError
91+
if errors.As(err, &apiErr) {
92+
return retryableErrorMap[apiErr.ErrorCode()]
93+
}
94+
return false
8795
}
8896

8997
// calculateWaitTime returns different time based on whether if

extension/entitystore/retryer_test.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@ import (
77
"testing"
88
"time"
99

10-
"github.com/aws/aws-sdk-go/aws/ec2metadata"
10+
"github.com/aws/aws-sdk-go-v2/feature/ec2/imds"
1111
"github.com/stretchr/testify/assert"
1212
"go.uber.org/zap"
1313

@@ -30,7 +30,7 @@ func TestRetryer_refreshLoop(t *testing.T) {
3030
name: "HappyPath_CorrectRefresh",
3131
fields: fields{
3232
metadataProvider: &mockMetadataProvider{
33-
InstanceIdentityDocument: &ec2metadata.EC2InstanceIdentityDocument{
33+
InstanceIdentityDocument: &imds.InstanceIdentityDocument{
3434
InstanceID: "i-123456789"},
3535
},
3636
iamRole: "original-role",

extension/entitystore/serviceprovider.go

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -10,7 +10,6 @@ import (
1010

1111
"go.uber.org/zap"
1212

13-
configaws "github.com/aws/amazon-cloudwatch-agent/cfg/aws"
1413
"github.com/aws/amazon-cloudwatch-agent/internal/ec2metadataprovider"
1514
"github.com/aws/amazon-cloudwatch-agent/plugins/processors/ec2tagger"
1615
"github.com/aws/amazon-cloudwatch-agent/translator/config"
@@ -317,7 +316,7 @@ func toLowerKeyMap(values []string) map[string]string {
317316
return set
318317
}
319318

320-
func newServiceProvider(mode string, region string, ec2Info *EC2Info, metadataProvider ec2metadataprovider.MetadataProvider, providerType ec2ProviderType, ec2Credential *configaws.CredentialConfig, done chan struct{}, logger *zap.Logger) serviceProviderInterface {
319+
func newServiceProvider(mode string, region string, ec2Info *EC2Info, metadataProvider ec2metadataprovider.MetadataProvider, done chan struct{}, logger *zap.Logger) serviceProviderInterface {
321320
return &serviceprovider{
322321
mode: mode,
323322
region: region,

extension/entitystore/serviceprovider_test.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,7 +8,7 @@ import (
88
"testing"
99
"time"
1010

11-
"github.com/aws/aws-sdk-go/aws/ec2metadata"
11+
"github.com/aws/aws-sdk-go-v2/feature/ec2/imds"
1212
"github.com/stretchr/testify/assert"
1313
"go.uber.org/zap"
1414

@@ -26,7 +26,7 @@ func Test_serviceprovider_startServiceProvider(t *testing.T) {
2626
{
2727
name: "HappyPath_AllServiceNames",
2828
metadataProvider: &mockMetadataProvider{
29-
InstanceIdentityDocument: &ec2metadata.EC2InstanceIdentityDocument{
29+
InstanceIdentityDocument: &imds.InstanceIdentityDocument{
3030
InstanceID: "i-123456789"},
3131
Tags: map[string]string{"service": "test-service"},
3232
},

0 commit comments

Comments
 (0)