diff --git a/pkg/recfilter/dedupe.go b/pkg/recfilter/dedupe.go index 8ff1cfd..9785bb3 100644 --- a/pkg/recfilter/dedupe.go +++ b/pkg/recfilter/dedupe.go @@ -116,13 +116,29 @@ func dedupeKey(resourceType, region, engine, deployment string) string { return fmt.Sprintf("%s|%s|%s|%s", resourceType, region, engine, deployment) } +const unknownElastiCacheEngine = "elasticache:*" + +func dedupeEngine(providerType common.ProviderType, service common.ServiceType, engine string) (string, bool) { + engine = common.NormalizeEngineName(engine) + if providerType != common.ProviderAWS || (service != common.ServiceCache && service != common.ServiceElastiCache) { + return engine, false + } + switch engine { + case "": + return unknownElastiCacheEngine, true + case "valkey": + engine = "redis" + } + return "elasticache:" + engine, true +} + // buildExistingCommitmentsMap builds a map of commitments by resource type, region, engine, and deployment. func buildExistingCommitmentsMap(commitments []common.Commitment, logf Logf) map[string]int { existingMap := make(map[string]int) for _rvc := range commitments { c := commitments[_rvc] - normalizedEngine := common.NormalizeEngineName(c.Engine) + normalizedEngine, _ := dedupeEngine(c.Provider, c.Service, c.Engine) normalizedDeployment := common.NormalizeDeploymentName(c.Deployment) key := dedupeKey(c.ResourceType, c.Region, normalizedEngine, normalizedDeployment) existingMap[key] += c.Count @@ -154,25 +170,33 @@ func adjustRecommendationsAgainstExisting(recs []common.Recommendation, existing // adjustSingleRecommendation adjusts a single recommendation based on existing commitments. func adjustSingleRecommendation(rec common.Recommendation, existingMap map[string]int, logf Logf) common.Recommendation { - engine := common.EngineFromDetails(rec.Details) + if rec.Count <= 0 { + return common.Recommendation{Count: 0} + } + engine, isElastiCache := dedupeEngine(rec.Provider, rec.Service, common.EngineFromDetails(rec.Details)) deployment := common.NormalizeDeploymentName(common.DeploymentFromDetails(rec.Details)) key := dedupeKey(rec.ResourceType, rec.Region, engine, deployment) - existingCount := existingMap[key] - - if existingCount >= rec.Count { - // All of this recommendation is covered by recent RIs. - // Return a zero-value Recommendation (Count=0) as a sentinel; the caller - // (adjustRecommendationsAgainstExisting) filters out recommendations with Count <= 0. - logf.printf(" [DuplicateChecker] SKIP %s: recent %d >= recommended %d", key, existingCount, rec.Count) - existingMap[key] -= rec.Count + keys := []string{key} + if isElastiCache && engine != unknownElastiCacheEngine { + keys = append(keys, dedupeKey(rec.ResourceType, rec.Region, unknownElastiCacheEngine, deployment)) + } + remaining := rec.Count + for _, candidate := range keys { + if available := existingMap[candidate]; available > 0 { + used := min(available, remaining) + remaining -= used + existingMap[candidate] -= used + } + } + + if remaining <= 0 { + logf.printf(" [DuplicateChecker] SKIP %s: recent commitments cover recommended %d", key, rec.Count) return common.Recommendation{Count: 0} } - // Partial or no coverage by recent RIs adjusted := rec - if existingCount > 0 { - adjusted.Count = rec.Count - existingCount - existingMap[key] = 0 + if remaining < rec.Count { + adjusted.Count = remaining logf.printf(" [DuplicateChecker] PARTIAL %s: adjusted count from %d to %d", key, rec.Count, adjusted.Count) } diff --git a/pkg/recfilter/dedupe_test.go b/pkg/recfilter/dedupe_test.go index 392dfe9..89ea71a 100644 --- a/pkg/recfilter/dedupe_test.go +++ b/pkg/recfilter/dedupe_test.go @@ -41,6 +41,132 @@ func (f *fakeServiceClient) GetValidResourceTypes(ctx context.Context) ([]string return nil, nil } +func TestElastiCacheEngineBudgets(t *testing.T) { + t.Parallel() + for _, tt := range []struct { + name string + engines []string + counts []int + recEngines []string + recCounts []int + wantCounts []int + }{ + {"redis covers valkey", []string{"ReDiS"}, []int{1}, []string{"VaLkEy"}, []int{1}, nil}, + {"valkey covers redis", []string{"valkey"}, []int{1}, []string{"redis"}, []int{1}, nil}, + {"redis not memcached", []string{"redis"}, []int{1}, []string{"memcached"}, []int{1}, []int{1}}, + {"valkey not memcached", []string{"valkey"}, []int{1}, []string{"memcached"}, []int{1}, []int{1}}, + {"unknown covers redis", []string{""}, []int{1}, []string{"redis"}, []int{1}, nil}, + {"unknown covers valkey", []string{""}, []int{1}, []string{"valkey"}, []int{1}, nil}, + {"unknown covers memcached", []string{""}, []int{1}, []string{"memcached"}, []int{1}, nil}, + {"combined partial", []string{"redis", ""}, []int{1, 1}, []string{"valkey"}, []int{3}, []int{1}}, + {"combined full", []string{"redis", ""}, []int{1, 1}, []string{"valkey"}, []int{2}, nil}, + {"wildcard consumed once", []string{""}, []int{1}, []string{"redis", "memcached"}, []int{1, 1}, []int{1}}, + {"exact before wildcard", []string{"redis", ""}, []int{1, 1}, []string{"valkey", "memcached"}, []int{1, 1}, nil}, + {"family consumed once", []string{"redis"}, []int{1}, []string{"valkey", "redis"}, []int{1, 1}, []int{1}}, + {"empty rec one", []string{""}, []int{1}, []string{""}, []int{3}, []int{2}}, + {"empty rec two", []string{""}, []int{2}, []string{""}, []int{3}, []int{1}}, + {"empty rec repeated", []string{""}, []int{2}, []string{"", ""}, []int{1, 2}, []int{1}}, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + client := &fakeServiceClient{} + for i, engine := range tt.engines { + client.commitments = append(client.commitments, common.Commitment{ + Provider: common.ProviderAWS, Service: common.ServiceCache, ResourceType: "cache.r6g.large", + Region: "us-east-1", Engine: engine, Count: tt.counts[i], State: common.CommitmentStateActive, StartDate: time.Now(), + }) + } + recs := make([]common.Recommendation, 0, len(tt.recEngines)) + for i, engine := range tt.recEngines { + recs = append(recs, common.Recommendation{Provider: common.ProviderAWS, Service: common.ServiceElastiCache, + ResourceType: "cache.r6g.large", Region: "us-east-1", Count: tt.recCounts[i], Details: &common.CacheDetails{Engine: engine}}) + } + passed, filtered, err := NewDuplicateChecker(0).AdjustRecommendationsForExisting(context.Background(), recs, client) + require.NoError(t, err) + require.Len(t, passed, len(tt.wantCounts)) + for i, rec := range passed { + assert.Equal(t, tt.wantCounts[i], rec.Count) + } + assert.Len(t, filtered, len(recs)-len(passed)) + }) + } +} + +func TestElastiCacheEngineScope(t *testing.T) { + t.Parallel() + for _, engine := range []string{"", "redis"} { + for _, other := range []struct { + provider common.ProviderType + service common.ServiceType + }{ + {common.ProviderAWS, common.ServiceMemoryDB}, {common.ProviderAWS, common.ServiceRDS}, + {common.ProviderAzure, common.ServiceCache}, + } { + for _, reverse := range []bool{false, true} { + t.Run(fmt.Sprintf("%s/%s/%s/reverse=%t", engine, other.provider, other.service, reverse), func(t *testing.T) { + t.Parallel() + cache := common.Recommendation{Provider: common.ProviderAWS, Service: common.ServiceCache, + ResourceType: "cache.r6g.large", Region: "us-east-1", Count: 2, Details: &common.CacheDetails{Engine: engine}} + foreign := cache + foreign.Provider, foreign.Service = other.provider, other.service + client := &fakeServiceClient{commitments: []common.Commitment{ + {Provider: cache.Provider, Service: cache.Service, ResourceType: cache.ResourceType, Region: cache.Region, + Engine: engine, Count: 1, State: common.CommitmentStateActive, StartDate: time.Now()}, + {Provider: other.provider, Service: other.service, ResourceType: cache.ResourceType, Region: cache.Region, + Engine: engine, Count: 2, State: common.CommitmentStateActive, StartDate: time.Now()}, + }} + recs := []common.Recommendation{cache, foreign} + if reverse { + recs[0], recs[1] = recs[1], recs[0] + } + passed, filtered, err := NewDuplicateChecker(0).AdjustRecommendationsForExisting(context.Background(), recs, client) + require.NoError(t, err) + require.Len(t, passed, 1) + assert.Equal(t, cache.Provider, passed[0].Provider) + assert.Equal(t, cache.Service, passed[0].Service) + assert.Equal(t, 1, passed[0].Count) + require.Len(t, filtered, 1) + assert.Equal(t, foreign.Provider, filtered[0].Provider) + assert.Equal(t, foreign.Service, filtered[0].Service) + }) + } + } + } +} + +func TestElastiCacheEngineBoundaries(t *testing.T) { + t.Parallel() + for _, tt := range []struct { + name string + provider common.ProviderType + service common.ServiceType + engine string + region string + resource string + }{ + {"different region", common.ProviderAWS, common.ServiceCache, "", "us-west-2", "cache.r6g.large"}, + {"different resource", common.ProviderAWS, common.ServiceCache, "", "us-east-1", "cache.r6g.xlarge"}, + {"memorydb unknown", common.ProviderAWS, common.ServiceMemoryDB, "", "us-east-1", "cache.r6g.large"}, + {"memorydb redis", common.ProviderAWS, common.ServiceMemoryDB, "redis", "us-east-1", "cache.r6g.large"}, + {"rds unknown", common.ProviderAWS, common.ServiceRDS, "", "us-east-1", "cache.r6g.large"}, + {"azure unknown", common.ProviderAzure, common.ServiceCache, "", "us-east-1", "cache.r6g.large"}, + {"azure redis", common.ProviderAzure, common.ServiceCache, "redis", "us-east-1", "cache.r6g.large"}, + } { + t.Run(tt.name, func(t *testing.T) { + t.Parallel() + client := &fakeServiceClient{commitments: []common.Commitment{{Provider: tt.provider, Service: tt.service, + ResourceType: tt.resource, Region: tt.region, Engine: tt.engine, Count: 1, + State: common.CommitmentStateActive, StartDate: time.Now()}}} + rec := common.Recommendation{Provider: tt.provider, Service: tt.service, + ResourceType: "cache.r6g.large", Region: "us-east-1", Count: 1, Details: &common.CacheDetails{Engine: "valkey"}} + passed, filtered, err := NewDuplicateChecker(0).AdjustRecommendationsForExisting(context.Background(), []common.Recommendation{rec}, client) + require.NoError(t, err) + assert.Equal(t, []common.Recommendation{rec}, passed) + assert.Empty(t, filtered) + }) + } +} + type customDedupeClient struct { fakeServiceClient filterErr error diff --git a/providers/aws/recommendations/parser_services_test.go b/providers/aws/recommendations/parser_services_test.go index b988da8..33e31ef 100644 --- a/providers/aws/recommendations/parser_services_test.go +++ b/providers/aws/recommendations/parser_services_test.go @@ -3,15 +3,87 @@ package recommendations import ( "context" "testing" + "time" "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/service/costexplorer/types" + ecSDK "github.com/aws/aws-sdk-go-v2/service/elasticache" + ecTypes "github.com/aws/aws-sdk-go-v2/service/elasticache/types" + mdbSDK "github.com/aws/aws-sdk-go-v2/service/memorydb" + mdbTypes "github.com/aws/aws-sdk-go-v2/service/memorydb/types" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/LeanerCloud/cloud-commitments-go/pkg/common" + "github.com/LeanerCloud/cloud-commitments-go/pkg/recfilter" + "github.com/LeanerCloud/cloud-commitments-go/providers/aws/services/elasticache" + "github.com/LeanerCloud/cloud-commitments-go/providers/aws/services/memorydb" ) +type cacheReservationAPI struct { + elasticache.API + engine *string +} + +func (a *cacheReservationAPI) DescribeReservedCacheNodes(context.Context, *ecSDK.DescribeReservedCacheNodesInput, ...func(*ecSDK.Options)) (*ecSDK.DescribeReservedCacheNodesOutput, error) { + return &ecSDK.DescribeReservedCacheNodesOutput{ReservedCacheNodes: []ecTypes.ReservedCacheNode{{ + ReservedCacheNodeId: aws.String("ri-new"), CacheNodeType: aws.String("cache.r6g.large"), + CacheNodeCount: aws.Int32(1), ProductDescription: a.engine, State: aws.String("active"), + Duration: aws.Int32(31536000), StartTime: aws.Time(time.Now().Add(-time.Hour)), + }}}, nil +} + +func TestParsedElastiCacheRecommendationDedupesSDKReservation(t *testing.T) { + t.Parallel() + for _, service := range []common.ServiceType{common.ServiceCache, common.ServiceElastiCache} { + for _, engine := range []*string{nil, aws.String(""), aws.String("redis")} { + t.Run(string(service)+"/"+aws.ToString(engine), func(t *testing.T) { + t.Parallel() + rec := common.Recommendation{Provider: common.ProviderAWS, Service: service, Count: 1} + err := (&Client{}).parseElastiCacheDetails(context.Background(), &rec, &types.ReservationPurchaseRecommendationDetail{ + InstanceDetails: &types.InstanceDetails{ElastiCacheInstanceDetails: &types.ElastiCacheInstanceDetails{ + NodeType: aws.String("cache.r6g.large"), Region: aws.String("us-east-1"), ProductDescription: aws.String("Valkey"), + }}, + }) + require.NoError(t, err) + client := elasticache.NewClient(aws.Config{Region: "us-east-1"}) + client.SetElastiCacheAPI(&cacheReservationAPI{engine: engine}) + passed, filtered, err := recfilter.NewDuplicateChecker(0).AdjustRecommendationsForExisting(context.Background(), []common.Recommendation{rec}, client) + require.NoError(t, err) + assert.Empty(t, passed) + assert.Len(t, filtered, 1) + assert.Equal(t, "Valkey", rec.Details.(*common.CacheDetails).Engine) + }) + } + } +} + +type memoryDBReservationAPI struct{ memorydb.API } + +func (*memoryDBReservationAPI) DescribeReservedNodes(context.Context, *mdbSDK.DescribeReservedNodesInput, ...func(*mdbSDK.Options)) (*mdbSDK.DescribeReservedNodesOutput, error) { + return &mdbSDK.DescribeReservedNodesOutput{ReservedNodes: []mdbTypes.ReservedNode{{ + ReservationId: aws.String("rn-new"), NodeType: aws.String("db.r6gd.xlarge"), NodeCount: 1, + State: aws.String("active"), Duration: 31536000, StartTime: aws.Time(time.Now().Add(-time.Hour)), + }}}, nil +} + +func TestParsedMemoryDBRecommendationDedupesSDKReservation(t *testing.T) { + t.Parallel() + rec := common.Recommendation{Provider: common.ProviderAWS, Service: common.ServiceMemoryDB, Count: 1} + err := (&Client{}).parseMemoryDBDetails(context.Background(), &rec, &types.ReservationPurchaseRecommendationDetail{ + InstanceDetails: &types.InstanceDetails{MemoryDBInstanceDetails: &types.MemoryDBInstanceDetails{ + NodeType: aws.String("db.r6gd.xlarge"), Region: aws.String("us-east-1"), + }}, + }) + require.NoError(t, err) + client := memorydb.NewClient(aws.Config{Region: "us-east-1"}) + client.SetMemoryDBAPI(&memoryDBReservationAPI{}) + passed, filtered, err := recfilter.NewDuplicateChecker(0).AdjustRecommendationsForExisting(context.Background(), []common.Recommendation{rec}, client) + require.NoError(t, err) + assert.Empty(t, passed) + assert.Len(t, filtered, 1) +} + func TestParseRDSDetails(t *testing.T) { client := &Client{} diff --git a/providers/aws/services/elasticache/client.go b/providers/aws/services/elasticache/client.go index 2584bf7..eca6a7d 100644 --- a/providers/aws/services/elasticache/client.go +++ b/providers/aws/services/elasticache/client.go @@ -90,6 +90,9 @@ func (c *Client) GetExistingCommitments(ctx context.Context) ([]common.Commitmen if duration == ThreeYearSeconds { termMonths = 36 } + if aws.ToString(node.ProductDescription) == "" { + log.Printf("WARNING: ElastiCache reservation %s has no engine; matching any cache engine during duplicate checks", aws.ToString(node.ReservedCacheNodeId)) + } commitment := common.Commitment{ Provider: common.ProviderAWS, diff --git a/providers/aws/services/elasticache/client_test.go b/providers/aws/services/elasticache/client_test.go index 8d1b894..a5c3430 100644 --- a/providers/aws/services/elasticache/client_test.go +++ b/providers/aws/services/elasticache/client_test.go @@ -1,8 +1,10 @@ package elasticache import ( + "bytes" "context" "fmt" + "log" "strings" "testing" "time" @@ -890,6 +892,8 @@ func TestClient_GetExistingCommitments_DedupesMatchingRecommendation(t *testing. client := &Client{client: mockClient, region: "us-east-1"} rec := common.Recommendation{ + Provider: common.ProviderAWS, + Service: common.ServiceCache, ResourceType: "cache.r6g.large", Region: "us-east-1", Count: 1, @@ -902,3 +906,24 @@ func TestClient_GetExistingCommitments_DedupesMatchingRecommendation(t *testing. assert.Empty(t, passed) assert.Len(t, filtered, 1) } + +func TestClient_GetExistingCommitments_WarnsOnMissingEngine(t *testing.T) { + var logs bytes.Buffer + previous := log.Writer() + log.SetOutput(&logs) + t.Cleanup(func() { log.SetOutput(previous) }) + for _, engine := range []*string{nil, aws.String("")} { + api := &MockElastiCacheClient{} + api.On("DescribeReservedCacheNodes", mock.Anything, mock.Anything).Return(&elasticache.DescribeReservedCacheNodesOutput{ + ReservedCacheNodes: []types.ReservedCacheNode{{ReservedCacheNodeId: aws.String("missing-engine-ri"), + State: aws.String("active"), ProductDescription: engine}}, + }, nil).Once() + commitments, err := (&Client{client: api, region: "us-east-1"}).GetExistingCommitments(context.Background()) + require.NoError(t, err) + require.Len(t, commitments, 1) + assert.Empty(t, commitments[0].Engine) + assert.Contains(t, logs.String(), "WARNING: ElastiCache reservation missing-engine-ri has no engine") + logs.Reset() + api.AssertExpectations(t) + } +} diff --git a/providers/aws/services/memorydb/client_test.go b/providers/aws/services/memorydb/client_test.go index 05b6edf..b17a68d 100644 --- a/providers/aws/services/memorydb/client_test.go +++ b/providers/aws/services/memorydb/client_test.go @@ -1086,38 +1086,3 @@ func TestFindOfferingID_InvalidTerm_ErrorsBeforeAPICall(t *testing.T) { } mockMDB.AssertNotCalled(t, "DescribeReservedNodesOfferings", mock.Anything, mock.Anything) } - -// Regression for cloud-commitments-go#66: the commitment must carry the same -// engine as the recommendation, or the dedupe key never matches and a node -// bought minutes ago is recommended and bought again. -func TestClient_GetExistingCommitments_DedupesMatchingRecommendation(t *testing.T) { - mockClient := &MockMemoryDBClient{} - t.Cleanup(func() { mockClient.AssertExpectations(t) }) - mockClient.On("DescribeReservedNodes", mock.Anything, mock.Anything). - Return(&memorydb.DescribeReservedNodesOutput{ - ReservedNodes: []types.ReservedNode{{ - ReservationId: aws.String("rn-new"), - NodeType: aws.String("db.r6gd.xlarge"), - NodeCount: 1, - State: aws.String("active"), - Duration: 31536000, - StartTime: aws.Time(time.Now().Add(-time.Hour)), - }}, - }, nil).Once() - - client := &Client{client: mockClient, region: "us-east-1"} - // The literal mirrors what parseMemoryDBDetails emits, so this fails if - // either side drifts from ReservedNodeEngine. - rec := common.Recommendation{ - ResourceType: "db.r6gd.xlarge", - Region: "us-east-1", - Count: 1, - Details: &common.CacheDetails{Engine: "redis", NodeType: "db.r6gd.xlarge"}, - } - - passed, filtered, err := recfilter.NewDuplicateChecker(0). - AdjustRecommendationsForExisting(context.Background(), []common.Recommendation{rec}, client) - require.NoError(t, err) - assert.Empty(t, passed) - assert.Len(t, filtered, 1) -}