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
52 changes: 38 additions & 14 deletions pkg/recfilter/dedupe.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
}

Expand Down
126 changes: 126 additions & 0 deletions pkg/recfilter/dedupe_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
72 changes: 72 additions & 0 deletions providers/aws/recommendations/parser_services_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{}

Expand Down
3 changes: 3 additions & 0 deletions providers/aws/services/elasticache/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
25 changes: 25 additions & 0 deletions providers/aws/services/elasticache/client_test.go
Original file line number Diff line number Diff line change
@@ -1,8 +1,10 @@
package elasticache

import (
"bytes"
"context"
"fmt"
"log"
"strings"
"testing"
"time"
Expand Down Expand Up @@ -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,
Expand All @@ -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)
}
}
35 changes: 0 additions & 35 deletions providers/aws/services/memorydb/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Loading