Skip to content
Merged
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
7 changes: 7 additions & 0 deletions extension/agenthealth/handler/stats/agent/flag.go
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,8 @@ type FlagSet interface {
SetValues(flags map[Flag]any)
// OnChange registers a callback that triggers on flag sets.
OnChange(callback func())
// Reset allows tests to clear the stored flags
Reset()
}

type flagSet struct {
Expand Down Expand Up @@ -180,6 +182,11 @@ func (p *flagSet) notify() {
}
}

func (p *flagSet) Reset() {
p.m.Clear()
p.notify()
}

func UsageFlags() FlagSet {
flagOnce.Do(func() {
flagSingleton = &flagSet{}
Expand Down
10 changes: 5 additions & 5 deletions extension/entitystore/ec2Info_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,22 +9,22 @@ import (
"testing"
"time"

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

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

var mockedInstanceIdentityDoc = &ec2metadata.EC2InstanceIdentityDocument{
var mockedInstanceIdentityDoc = &imds.InstanceIdentityDocument{
InstanceID: "i-01d2417c27a396e44",
AccountID: "874389809020",
Region: "us-east-1",
InstanceType: "m5ad.large",
ImageID: "ami-09edd32d9b0990d49",
}

var mockedInstanceIdentityDocWithLargeInstanceId = &ec2metadata.EC2InstanceIdentityDocument{
var mockedInstanceIdentityDocWithLargeInstanceID = &imds.InstanceIdentityDocument{
InstanceID: "i-01d2417c27a396e44394824728",
AccountID: "874389809020",
Region: "us-east-1",
Expand Down Expand Up @@ -60,12 +60,12 @@ func TestSetInstanceIDAccountID(t *testing.T) {
{
name: "InstanceId too large",
args: args{
metadataProvider: &mockMetadataProvider{InstanceIdentityDocument: mockedInstanceIdentityDocWithLargeInstanceId},
metadataProvider: &mockMetadataProvider{InstanceIdentityDocument: mockedInstanceIdentityDocWithLargeInstanceID},
},
wantErr: false,
want: EC2Info{
InstanceID: "",
AccountID: mockedInstanceIdentityDocWithLargeInstanceId.AccountID,
AccountID: mockedInstanceIdentityDocWithLargeInstanceID.AccountID,
},
},
}
Expand Down
53 changes: 13 additions & 40 deletions extension/entitystore/extension.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,20 +7,17 @@ import (
"context"
"time"

"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/service/cloudwatchlogs/types"
"github.com/aws/aws-sdk-go/aws"
"github.com/aws/aws-sdk-go/aws/client"
"github.com/aws/aws-sdk-go/service/ec2"
"github.com/aws/aws-sdk-go/service/ec2/ec2iface"
"github.com/jellydator/ttlcache/v3"
"go.opentelemetry.io/collector/component"
"go.opentelemetry.io/collector/extension"
"go.uber.org/atomic"
"go.uber.org/zap"

configaws "github.com/aws/amazon-cloudwatch-agent/cfg/aws"
configaws "github.com/aws/amazon-cloudwatch-agent/cfg/aws/v2"
"github.com/aws/amazon-cloudwatch-agent/internal/ec2metadataprovider"
"github.com/aws/amazon-cloudwatch-agent/internal/retryer"
"github.com/aws/amazon-cloudwatch-agent/internal/retryer/v2"
"github.com/aws/amazon-cloudwatch-agent/plugins/processors/awsentity/entityattributes"
"github.com/aws/amazon-cloudwatch-agent/translator/config"
)
Expand All @@ -35,8 +32,6 @@ const (
podTerminationCheckInterval = 5 * time.Minute
)

type ec2ProviderType func(string, *configaws.CredentialConfig) ec2iface.EC2API

type serviceProviderInterface interface {
startServiceProvider()
addEntryForLogFile(LogFileGlob, ServiceAttribute)
Expand Down Expand Up @@ -69,31 +64,23 @@ type EntityStore struct {
// that we can attach to the entity
serviceprovider serviceProviderInterface

// nativeCredential stores the credential config for agent's native
// component such as LogAgent
nativeCredential client.ConfigProvider

metadataprovider ec2metadataprovider.MetadataProvider

podTerminationCheckInterval time.Duration
}

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

func (e *EntityStore) Start(ctx context.Context, host component.Host) error {
func (e *EntityStore) Start(ctx context.Context, _ component.Host) error {
// Get IMDS client and EC2 API client which requires region for authentication
// These will be passed down to any object that requires access to IMDS or EC2
// API client so we have single source of truth for credential
e.done = make(chan struct{})
e.metadataprovider = getMetaDataProvider()
e.metadataprovider = getMetaDataProvider(ctx)
e.mode = e.config.Mode
e.kubernetesMode = e.config.KubernetesMode
e.podTerminationCheckInterval = podTerminationCheckInterval
ec2CredentialConfig := &configaws.CredentialConfig{
Profile: e.config.Profile,
Filename: e.config.Filename,
}
e.serviceprovider = newServiceProvider(e.mode, e.config.Region, &e.ec2Info, e.metadataprovider, getEC2Provider, ec2CredentialConfig, e.done, e.logger)
e.serviceprovider = newServiceProvider(e.mode, e.config.Region, &e.ec2Info, e.metadataprovider, e.done, e.logger)
switch e.mode {
case config.ModeEC2:
e.ec2Info = *newEC2Info(e.metadataprovider, e.done, e.config.Region, e.logger)
Expand Down Expand Up @@ -138,14 +125,6 @@ func (e *EntityStore) EC2Info() EC2Info {
return e.ec2Info
}

func (e *EntityStore) SetNativeCredential(client client.ConfigProvider) {
e.nativeCredential = client
}

func (e *EntityStore) NativeCredentialExists() bool {
return e.nativeCredential != nil
}

// CreateLogFileEntity creates the entity for log events that are being uploaded from a log file in the environment.
func (e *EntityStore) CreateLogFileEntity(logFileGlob LogFileGlob, logGroupName LogGroupName) *types.Entity {
if e.serviceprovider == nil {
Expand Down Expand Up @@ -260,19 +239,13 @@ func (e *EntityStore) createServiceKeyAttributes(serviceAttr ServiceAttribute) m
return serviceKeyAttr
}

var getMetaDataProvider = func() ec2metadataprovider.MetadataProvider {
mdCredentialConfig := &configaws.CredentialConfig{}
return ec2metadataprovider.NewMetadataProvider(mdCredentialConfig.Credentials(), retryer.GetDefaultRetryNumber())
}

var getEC2Provider = func(region string, ec2CredentialConfig *configaws.CredentialConfig) ec2iface.EC2API {
ec2CredentialConfig.Region = region
return ec2.New(
ec2CredentialConfig.Credentials(),
&aws.Config{
LogLevel: configaws.SDKLogLevel(),
Logger: configaws.SDKLogger{},
})
var getMetaDataProvider = func(ctx context.Context) ec2metadataprovider.MetadataProvider {
mdCredentialConfig := &configaws.CredentialsConfig{}
cfg, err := mdCredentialConfig.LoadConfig(ctx)
if err != nil {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

should log the error?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It's already logged in config/aws/v2/credentials.go.

cfg = aws.Config{}
}
return ec2metadataprovider.NewMetadataProvider(cfg, retryer.GetDefaultRetryNumber())
}

func addNonEmptyToMap(m map[string]string, key, value string) {
Expand Down
41 changes: 19 additions & 22 deletions extension/entitystore/extension_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -12,9 +12,8 @@ import (
"testing"
"time"

"github.com/aws/aws-sdk-go-v2/feature/ec2/imds"
"github.com/aws/aws-sdk-go-v2/service/cloudwatchlogs/types"
"github.com/aws/aws-sdk-go/aws/ec2metadata"
"github.com/aws/aws-sdk-go/aws/session"
"github.com/jellydator/ttlcache/v3"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/mock"
Expand Down Expand Up @@ -73,57 +72,59 @@ func (s *mockServiceProvider) setAutoScalingGroup(asg string) {
}

type mockMetadataProvider struct {
InstanceIdentityDocument *ec2metadata.EC2InstanceIdentityDocument
InstanceIdentityDocument *imds.InstanceIdentityDocument
Tags map[string]string
InstanceTagError bool
}

func mockMetadataProviderFunc() ec2metadataprovider.MetadataProvider {
var _ ec2metadataprovider.MetadataProvider = (*mockMetadataProvider)(nil)

func mockMetadataProviderFunc(context.Context) ec2metadataprovider.MetadataProvider {
return &mockMetadataProvider{
Tags: map[string]string{
"aws:autoscaling:groupName": "ASG-1",
},
InstanceIdentityDocument: &ec2metadata.EC2InstanceIdentityDocument{
InstanceIdentityDocument: &imds.InstanceIdentityDocument{
InstanceID: "i-123456789",
},
}
}

func mockMetadataProviderWithAccountId(accountId string) *mockMetadataProvider {
return &mockMetadataProvider{
InstanceIdentityDocument: &ec2metadata.EC2InstanceIdentityDocument{
InstanceIdentityDocument: &imds.InstanceIdentityDocument{
AccountID: accountId,
},
}
}

func (m *mockMetadataProvider) Get(ctx context.Context) (ec2metadata.EC2InstanceIdentityDocument, error) {
func (m *mockMetadataProvider) Get(context.Context) (imds.InstanceIdentityDocument, error) {
if m.InstanceIdentityDocument != nil {
return *m.InstanceIdentityDocument, nil
}
return ec2metadata.EC2InstanceIdentityDocument{}, errors.New("No instance identity document")
return imds.InstanceIdentityDocument{}, errors.New("no instance identity document")
}

func (m *mockMetadataProvider) Hostname(ctx context.Context) (string, error) {
func (m *mockMetadataProvider) Hostname(context.Context) (string, error) {
return "MockHostName", nil
}

func (m *mockMetadataProvider) InstanceID(ctx context.Context) (string, error) {
func (m *mockMetadataProvider) InstanceID(context.Context) (string, error) {
return "MockInstanceID", nil
}

func (m *mockMetadataProvider) InstanceTags(_ context.Context) ([]string, error) {
func (m *mockMetadataProvider) InstanceTags(context.Context) ([]string, error) {
if m.InstanceTagError {
return nil, errors.New("an error occurred for instance tag retrieval")
}
return maps.Keys(m.Tags), nil
}

func (m *mockMetadataProvider) ClientIAMRole(ctx context.Context) (string, error) {
func (m *mockMetadataProvider) ClientIAMRole(context.Context) (string, error) {
return "TestRole", nil
}

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

entity := e.CreateLogFileEntity(glob, group)
Expand Down Expand Up @@ -364,9 +364,8 @@ func TestEntityStore_createLogFileRID_ServiceProviderIsEmpty(t *testing.T) {
glob := LogFileGlob("glob")
group := LogGroupName("group")
e := EntityStore{
mode: config.ModeEC2,
ec2Info: EC2Info{InstanceID: instanceId},
nativeCredential: &session.Session{},
mode: config.ModeEC2,
ec2Info: EC2Info{InstanceID: instanceId},
}

entity := e.CreateLogFileEntity(glob, group)
Expand Down Expand Up @@ -548,7 +547,6 @@ func TestEntityStore_GetMetricServiceNameSource(t *testing.T) {
ec2Info: EC2Info{InstanceID: instanceId},
serviceprovider: sp,
metadataprovider: mockMetadataProviderWithAccountId(accountId),
nativeCredential: &session.Session{},
}

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

serviceName, serviceNameSource := e.GetMetricServiceNameAndSource()
Expand Down
18 changes: 13 additions & 5 deletions extension/entitystore/retryer.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,16 +4,16 @@
package entitystore

import (
"errors"
"math/rand"
"time"

"github.com/aws/aws-sdk-go/aws/awserr"
"github.com/aws/smithy-go"
"go.uber.org/zap"
)

const (
RequestLimitExceeded = "RequestLimitExceeded"
infRetry = -1
infRetry = -1
)

var (
Expand Down Expand Up @@ -57,7 +57,7 @@ func (r *Retryer) refreshLoop(updateFunc func() error) int {
err := updateFunc()
if err == nil && r.oneTime {
return retry
} else if awsErr, ok := err.(awserr.Error); ok && !r.retryAnyError && !retryableErrorMap[awsErr.Code()] {
} else if !r.retryAnyError && !isRetryableError(err) {
return retry
}

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

}
return retry
}

// isRetryableError checks if the error is a retryable API error code recognized by the extension.
func isRetryableError(err error) bool {
var apiErr smithy.APIError
if errors.As(err, &apiErr) {
return retryableErrorMap[apiErr.ErrorCode()]
}
return false
}

// calculateWaitTime returns different time based on whether if
Expand Down
4 changes: 2 additions & 2 deletions extension/entitystore/retryer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ import (
"testing"
"time"

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

Expand All @@ -30,7 +30,7 @@ func TestRetryer_refreshLoop(t *testing.T) {
name: "HappyPath_CorrectRefresh",
fields: fields{
metadataProvider: &mockMetadataProvider{
InstanceIdentityDocument: &ec2metadata.EC2InstanceIdentityDocument{
InstanceIdentityDocument: &imds.InstanceIdentityDocument{
InstanceID: "i-123456789"},
},
iamRole: "original-role",
Expand Down
3 changes: 1 addition & 2 deletions extension/entitystore/serviceprovider.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,7 +10,6 @@ import (

"go.uber.org/zap"

configaws "github.com/aws/amazon-cloudwatch-agent/cfg/aws"
"github.com/aws/amazon-cloudwatch-agent/internal/ec2metadataprovider"
"github.com/aws/amazon-cloudwatch-agent/plugins/processors/ec2tagger"
"github.com/aws/amazon-cloudwatch-agent/translator/config"
Expand Down Expand Up @@ -317,7 +316,7 @@ func toLowerKeyMap(values []string) map[string]string {
return set
}

func newServiceProvider(mode string, region string, ec2Info *EC2Info, metadataProvider ec2metadataprovider.MetadataProvider, providerType ec2ProviderType, ec2Credential *configaws.CredentialConfig, done chan struct{}, logger *zap.Logger) serviceProviderInterface {
func newServiceProvider(mode string, region string, ec2Info *EC2Info, metadataProvider ec2metadataprovider.MetadataProvider, done chan struct{}, logger *zap.Logger) serviceProviderInterface {
return &serviceprovider{
mode: mode,
region: region,
Expand Down
4 changes: 2 additions & 2 deletions extension/entitystore/serviceprovider_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ import (
"testing"
"time"

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

Expand All @@ -26,7 +26,7 @@ func Test_serviceprovider_startServiceProvider(t *testing.T) {
{
name: "HappyPath_AllServiceNames",
metadataProvider: &mockMetadataProvider{
InstanceIdentityDocument: &ec2metadata.EC2InstanceIdentityDocument{
InstanceIdentityDocument: &imds.InstanceIdentityDocument{
InstanceID: "i-123456789"},
Tags: map[string]string{"service": "test-service"},
},
Expand Down
Loading
Loading