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
46 changes: 45 additions & 1 deletion providers/aws/recommendations/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import (
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/service/costexplorer"
"github.com/aws/aws-sdk-go-v2/service/costexplorer/types"
awsec2 "github.com/aws/aws-sdk-go-v2/service/ec2"
"golang.org/x/sync/errgroup"

"github.com/LeanerCloud/CUDly/pkg/common"
Expand All @@ -34,6 +35,20 @@ type Client struct {
costExplorerClient CostExplorerAPI
region string
rateLimiter *RateLimiter

// ec2API is the EC2 client used to build the DescribeInstanceTypes paginator.
// Populated by NewClient from aws.Config; nil when created via NewClientWithAPI.
ec2API DescribeInstanceTypesAPI

// instanceTypePagerFactory creates a new InstanceTypePager on demand.
// Set by NewClient to wrap ec2API; overridable via SetInstanceTypePagerFactory
// for hermetic tests. When nil, instanceTypeLookup returns (0,0).
instanceTypePagerFactory func() InstanceTypePager

// skuCatalog caches the per-instance-type vCPU/memory catalogue, fetched
// lazily once per Client lifetime via sync.Once (one DescribeInstanceTypes
// fan-out per scheduler tick).
skuCatalog skuCatalog
}

// NewClient creates a new recommendations client
Expand All @@ -43,10 +58,17 @@ func NewClient(cfg aws.Config) *Client {
ceConfig.Region = "us-east-1"
ceConfig.BaseEndpoint = aws.String("https://ce.us-east-1.amazonaws.com")

ec2Client := awsec2.NewFromConfig(cfg)
return &Client{
costExplorerClient: costexplorer.NewFromConfig(ceConfig),
region: cfg.Region,
rateLimiter: NewRateLimiter(),
ec2API: ec2Client,
// Factory wraps the EC2 client so the paginator is created lazily
// on the first EC2 recommendation parse (not at construction time).
instanceTypePagerFactory: func() InstanceTypePager {
return awsec2.NewDescribeInstanceTypesPaginator(ec2Client, &awsec2.DescribeInstanceTypesInput{})
},
}
}

Expand All @@ -56,7 +78,29 @@ func NewClientWithAPI(api CostExplorerAPI, region string) *Client {
costExplorerClient: api,
region: region,
rateLimiter: NewRateLimiter(),
// ec2API left nil: instanceTypeLookup falls back to VCPU=0/MemoryGB=0
// unless the caller sets instanceTypePagerFactory.
}
}

// SetInstanceTypePagerFactory injects a pager factory for the instance-type
// SKU catalogue. Must be called before the first GetRecommendations call.
// Intended for tests that need to verify the one-fetch-per-lifetime invariant
// without hitting AWS.
func (c *Client) SetInstanceTypePagerFactory(f func() InstanceTypePager) {
c.instanceTypePagerFactory = f
}

// instanceTypeLookup returns the cached SKU entry for instanceType.
// On the first call the catalogue is built by calling the pager factory.
// ok=false when no factory is configured, the catalogue fetch failed, or
// the instance type was not in the catalogue — the caller falls back to
// VCPU=0/MemoryGB=0 (graceful-degradation contract from Azure PR #810).
func (c *Client) instanceTypeLookup(ctx context.Context, instanceType string) (instanceTypeSKUEntry, bool) {
if c.instanceTypePagerFactory == nil {
return instanceTypeSKUEntry{}, false
}
return c.skuCatalog.lookup(ctx, instanceType, c.instanceTypePagerFactory)
}

// GetRecommendations fetches Reserved Instance recommendations for any service
Expand All @@ -83,7 +127,7 @@ func (c *Client) GetRecommendations(ctx context.Context, params common.Recommend
return nil, err
}

return c.parseRecommendations(allRecs, params)
return c.parseRecommendations(ctx, allRecs, params)
}

// fetchRIAllPages paginates over all pages of RI recommendations for a single
Expand Down
15 changes: 8 additions & 7 deletions providers/aws/recommendations/parser_ri.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package recommendations

import (
"context"
"fmt"
"log"
"math"
Expand All @@ -14,12 +15,12 @@ import (
)

// parseRecommendations converts AWS recommendations to common.Recommendation format
func (c *Client) parseRecommendations(awsRecs []types.ReservationPurchaseRecommendation, params common.RecommendationParams) ([]common.Recommendation, error) {
func (c *Client) parseRecommendations(ctx context.Context, awsRecs []types.ReservationPurchaseRecommendation, params common.RecommendationParams) ([]common.Recommendation, error) {
var recommendations []common.Recommendation

for _, awsRec := range awsRecs {
for i, details := range awsRec.RecommendationDetails {
rec, err := c.parseRecommendationDetail(&details, params)
rec, err := c.parseRecommendationDetail(ctx, &details, params)
if err != nil {
fmt.Printf("Warning: Failed to parse recommendation detail %d: %v\n", i, err)
continue
Expand All @@ -35,7 +36,7 @@ func (c *Client) parseRecommendations(awsRecs []types.ReservationPurchaseRecomme
}

// parseRecommendationDetail converts a single AWS recommendation detail
func (c *Client) parseRecommendationDetail(details *types.ReservationPurchaseRecommendationDetail, params common.RecommendationParams) (*common.Recommendation, error) {
func (c *Client) parseRecommendationDetail(ctx context.Context, details *types.ReservationPurchaseRecommendationDetail, params common.RecommendationParams) (*common.Recommendation, error) {
rec := &common.Recommendation{
Provider: common.ProviderAWS,
Service: params.Service,
Expand Down Expand Up @@ -74,7 +75,7 @@ func (c *Client) parseRecommendationDetail(details *types.ReservationPurchaseRec
c.parseRIUtilizationSignals(rec, details)

// Parse service-specific details
if err := c.parseServiceSpecificDetails(rec, details, params.Service); err != nil {
if err := c.parseServiceSpecificDetails(ctx, rec, details, params.Service); err != nil {
return nil, err
}

Expand Down Expand Up @@ -178,10 +179,10 @@ func (c *Client) parseAWSCostDetails(rec *common.Recommendation, details *types.
}

// serviceParserFunc defines the signature for service-specific parsers
type serviceParserFunc func(*common.Recommendation, *types.ReservationPurchaseRecommendationDetail) error
type serviceParserFunc func(context.Context, *common.Recommendation, *types.ReservationPurchaseRecommendationDetail) error

// parseServiceSpecificDetails routes to the appropriate service parser
func (c *Client) parseServiceSpecificDetails(rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail, service common.ServiceType) error {
func (c *Client) parseServiceSpecificDetails(ctx context.Context, rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail, service common.ServiceType) error {
// Map of service types to their parser functions
serviceParsers := map[common.ServiceType]serviceParserFunc{
common.ServiceRDS: c.parseRDSDetails,
Expand All @@ -202,5 +203,5 @@ func (c *Client) parseServiceSpecificDetails(rec *common.Recommendation, details
return fmt.Errorf("unsupported service: %s", service)
}

return parser(rec, details)
return parser(ctx, rec, details)
}
13 changes: 7 additions & 6 deletions providers/aws/recommendations/parser_ri_test.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package recommendations

import (
"context"
"testing"

"github.com/aws/aws-sdk-go-v2/aws"
Expand Down Expand Up @@ -184,7 +185,7 @@ func TestParseRecommendationDetail_UnsupportedService(t *testing.T) {
LookbackPeriod: "7d",
}

rec, err := client.parseRecommendationDetail(details, params)
rec, err := client.parseRecommendationDetail(context.Background(), details, params)

assert.Error(t, err)
assert.Nil(t, rec)
Expand All @@ -207,7 +208,7 @@ func TestParseRecommendationDetail_MissingQuantity(t *testing.T) {
LookbackPeriod: "7d",
}

rec, err := client.parseRecommendationDetail(details, params)
rec, err := client.parseRecommendationDetail(context.Background(), details, params)

assert.Error(t, err)
assert.Nil(t, rec)
Expand Down Expand Up @@ -241,7 +242,7 @@ func TestParseRecommendationDetail_WithAccountAndCosts(t *testing.T) {
LookbackPeriod: "7d",
}

rec, err := client.parseRecommendationDetail(details, params)
rec, err := client.parseRecommendationDetail(context.Background(), details, params)

require.NoError(t, err)
require.NotNil(t, rec)
Expand Down Expand Up @@ -300,7 +301,7 @@ func TestParseRecommendations(t *testing.T) {
LookbackPeriod: "7d",
}

recs, err := client.parseRecommendations(awsRecs, params)
recs, err := client.parseRecommendations(context.Background(), awsRecs, params)

require.NoError(t, err)
assert.Len(t, recs, 2)
Expand Down Expand Up @@ -371,7 +372,7 @@ func TestParseRecommendations_SkipsInvalidDetails(t *testing.T) {
LookbackPeriod: "7d",
}

recs, err := client.parseRecommendations(awsRecs, params)
recs, err := client.parseRecommendations(context.Background(), awsRecs, params)

require.NoError(t, err)
// Should have 2 valid recommendations, skipping the invalid one
Expand All @@ -390,7 +391,7 @@ func TestParseRecommendations_EmptyInput(t *testing.T) {
LookbackPeriod: "7d",
}

recs, err := client.parseRecommendations([]types.ReservationPurchaseRecommendation{}, params)
recs, err := client.parseRecommendations(context.Background(), []types.ReservationPurchaseRecommendation{}, params)

require.NoError(t, err)
assert.Empty(t, recs)
Expand Down
76 changes: 52 additions & 24 deletions providers/aws/recommendations/parser_services.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package recommendations

import (
"context"
"fmt"
"strings"

Expand All @@ -11,7 +12,7 @@ import (
)

// parseRDSDetails extracts RDS-specific details
func (c *Client) parseRDSDetails(rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail) error {
func (c *Client) parseRDSDetails(_ context.Context, rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail) error {
if details.InstanceDetails == nil || details.InstanceDetails.RDSInstanceDetails == nil {
return fmt.Errorf("RDS instance details not found")
}
Expand Down Expand Up @@ -43,7 +44,7 @@ func (c *Client) parseRDSDetails(rec *common.Recommendation, details *types.Rese
}

// parseElastiCacheDetails extracts ElastiCache-specific details
func (c *Client) parseElastiCacheDetails(rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail) error {
func (c *Client) parseElastiCacheDetails(_ context.Context, rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail) error {
if details.InstanceDetails == nil || details.InstanceDetails.ElastiCacheInstanceDetails == nil {
return fmt.Errorf("ElastiCache instance details not found")
}
Expand All @@ -66,8 +67,48 @@ func (c *Client) parseElastiCacheDetails(rec *common.Recommendation, details *ty
return nil
}

// parseEC2Details extracts EC2-specific details
func (c *Client) parseEC2Details(rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail) error {
// resolveEC2Tenancy maps a Cost Explorer tenancy value to the EC2 RI API
// tenancy string. CE uses "shared" for the default tenancy; "dedicated" maps
// directly. Any nil or unrecognised value is treated as default.
func resolveEC2Tenancy(tenancy *string) string {
if tenancy != nil && *tenancy == "dedicated" {
return string(ec2types.TenancyDedicated)
}
return string(ec2types.TenancyDefault)
}

// resolveEC2Scope maps a Cost Explorer availability zone value to the EC2 RI
// API scope string. A non-empty AZ means AZ scope; otherwise region scope.
func resolveEC2Scope(az *string) string {
if az != nil && *az != "" {
return string(ec2types.ScopeAvailabilityZone)
}
return string(ec2types.ScopeRegional)
}

// enrichFromCatalogue populates VCPU and MemoryGB on ec2Info from the
// lazily-cached DescribeInstanceTypes catalogue. Non-fatal on cache miss.
func (c *Client) enrichFromCatalogue(ctx context.Context, ec2Info *common.ComputeDetails) {
if ec2Info.InstanceType == "" {
return
}
entry, ok := c.instanceTypeLookup(ctx, ec2Info.InstanceType)
if !ok {
return
}
if entry.vCPUs > 0 {
ec2Info.VCPU = entry.vCPUs
}
if entry.memoryGB > 0 {
ec2Info.MemoryGB = entry.memoryGB
}
}

// parseEC2Details extracts EC2-specific details and enriches the rec with
// vCPU and memory from the lazily-cached DescribeInstanceTypes catalogue.
// If the catalogue fetch failed or the instance type is not found, VCPU
// and MemoryGB remain 0 (the omitempty JSON tags hide them from payloads).
func (c *Client) parseEC2Details(ctx context.Context, rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail) error {
if details.InstanceDetails == nil || details.InstanceDetails.EC2InstanceDetails == nil {
return fmt.Errorf("EC2 instance details not found")
}
Expand All @@ -85,29 +126,16 @@ func (c *Client) parseEC2Details(rec *common.Recommendation, details *types.Rese
if ec2Details.Region != nil {
rec.Region = normalizeRegionName(*ec2Details.Region)
}
// Tenancy: CE returns "shared" for default tenancy; the EC2 RI filter API
// expects "default" (types.TenancyDefault). CE "dedicated" maps directly.
// Any nil or unrecognised value is treated as default.
if ec2Details.Tenancy != nil && *ec2Details.Tenancy == "dedicated" {
ec2Info.Tenancy = string(ec2types.TenancyDedicated)
} else {
ec2Info.Tenancy = string(ec2types.TenancyDefault)
}

// Scope: the EC2 RI filter API expects "Region" (types.ScopeRegional) or
// "Availability Zone" (types.ScopeAvailabilityZone) - not lowercase/hyphenated.
if ec2Details.AvailabilityZone != nil && *ec2Details.AvailabilityZone != "" {
ec2Info.Scope = string(ec2types.ScopeAvailabilityZone)
} else {
ec2Info.Scope = string(ec2types.ScopeRegional)
}
ec2Info.Tenancy = resolveEC2Tenancy(ec2Details.Tenancy)
ec2Info.Scope = resolveEC2Scope(ec2Details.AvailabilityZone)
c.enrichFromCatalogue(ctx, ec2Info)

rec.Details = ec2Info
return nil
}

// parseOpenSearchDetails extracts OpenSearch-specific details
func (c *Client) parseOpenSearchDetails(rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail) error {
func (c *Client) parseOpenSearchDetails(_ context.Context, rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail) error {
if details.InstanceDetails == nil || details.InstanceDetails.ESInstanceDetails == nil {
return fmt.Errorf("OpenSearch/Elasticsearch instance details not found")
}
Expand All @@ -134,7 +162,7 @@ func (c *Client) parseOpenSearchDetails(rec *common.Recommendation, details *typ
}

// parseRedshiftDetails extracts Redshift-specific details
func (c *Client) parseRedshiftDetails(rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail) error {
func (c *Client) parseRedshiftDetails(_ context.Context, rec *common.Recommendation, details *types.ReservationPurchaseRecommendationDetail) error {
if details.InstanceDetails == nil || details.InstanceDetails.RedshiftInstanceDetails == nil {
return fmt.Errorf("Redshift instance details not found")
}
Expand Down Expand Up @@ -172,9 +200,9 @@ func (c *Client) parseRedshiftDetails(rec *common.Recommendation, details *types
// If rec.ResourceType is empty, the function returns an error so the
// recommendation is skipped loudly (logged by parseRecommendations) rather
// than silently substituting a wrong default instance type.
func (c *Client) parseMemoryDBDetails(rec *common.Recommendation, _ *types.ReservationPurchaseRecommendationDetail) error {
func (c *Client) parseMemoryDBDetails(_ context.Context, rec *common.Recommendation, _ *types.ReservationPurchaseRecommendationDetail) error {
if rec.ResourceType == "" {
return fmt.Errorf("MemoryDB recommendation has no ResourceType; cannot determine offering — Cost Explorer did not populate instance details")
return fmt.Errorf("MemoryDB recommendation has no ResourceType; cannot determine offering - Cost Explorer did not populate instance details")
}
rec.Details = &common.CacheDetails{
Engine: "redis",
Expand Down
Loading
Loading