From 9a40c12dc7881439f011bd6bdacdf6f1c565bf36 Mon Sep 17 00:00:00 2001 From: Cristian Magherusan-Stanciu Date: Fri, 2 Oct 2026 20:43:58 +0200 Subject: [PATCH] fix(scheduler): retain incomplete AWS recommendation sweeps Preserve valid rows and withhold stale-row eviction for incomplete AWS accounts. Query explicit services when retrying with configured recommendation parameters. Refs LeanerCloud/cloud-commitments-go#54 --- go.mod | 4 +- go.sum | 8 +- .../scheduler/partial_sweep_eviction_test.go | 71 +++++ ...mendation_completeness_integration_test.go | 246 ++++++++++++++++++ internal/scheduler/scheduler.go | 44 ++-- internal/scheduler/scheduler_test.go | 15 +- 6 files changed, 360 insertions(+), 28 deletions(-) create mode 100644 internal/scheduler/recommendation_completeness_integration_test.go diff --git a/go.mod b/go.mod index d778fc2c..0c485247 100644 --- a/go.mod +++ b/go.mod @@ -84,8 +84,8 @@ require ( github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/reservations/armreservations v1.1.0 github.com/Azure/azure-sdk-for-go/sdk/security/keyvault/azkeys v1.4.0 github.com/Azure/azure-sdk-for-go/sdk/security/keyvault/azsecrets v1.4.0 - github.com/LeanerCloud/cloud-commitments-go/pkg v0.0.0-20260928214714-ce9513612901 - github.com/LeanerCloud/cloud-commitments-go/providers/aws v0.0.0-20260928214714-ce9513612901 + github.com/LeanerCloud/cloud-commitments-go/pkg v0.0.0-20260929105827-b3b4cb5e3d80 + github.com/LeanerCloud/cloud-commitments-go/providers/aws v0.0.0-20261002152209-006ef5c8d0a2 github.com/LeanerCloud/cloud-commitments-go/providers/azure v0.0.0-20260928214714-ce9513612901 github.com/LeanerCloud/cloud-commitments-go/providers/gcp v0.0.0-20260928214714-ce9513612901 github.com/aws/aws-lambda-go v1.47.0 diff --git a/go.sum b/go.sum index f5cb3094..0f9acf91 100644 --- a/go.sum +++ b/go.sum @@ -92,10 +92,10 @@ github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/cloudmock v0 github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/cloudmock v0.54.0/go.mod h1:vB2GH9GAYYJTO3mEn8oYwzEdhlayZIdQz6zdzgUIRvA= github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/resourcemapping v0.54.0 h1:s0WlVbf9qpvkh1c/uDAPElam0WrL7fHRIidgZJ7UqZI= github.com/GoogleCloudPlatform/opentelemetry-operations-go/internal/resourcemapping v0.54.0/go.mod h1:Mf6O40IAyB9zR/1J8nGDDPirZQQPbYJni8Yisy7NTMc= -github.com/LeanerCloud/cloud-commitments-go/pkg v0.0.0-20260928214714-ce9513612901 h1:JWQwkwshjBoEuoIuFK1KHh70yZDYQTrm5GRloD6ojjU= -github.com/LeanerCloud/cloud-commitments-go/pkg v0.0.0-20260928214714-ce9513612901/go.mod h1:ApWBliDXe099f3oDXBz41K/I9v4bHvn1dG/BGoRmHlw= -github.com/LeanerCloud/cloud-commitments-go/providers/aws v0.0.0-20260928214714-ce9513612901 h1:18MxQd7+ZQ+Z2uNbZfpIJQ00JWdSdTVSmVWBzg4R0DA= -github.com/LeanerCloud/cloud-commitments-go/providers/aws v0.0.0-20260928214714-ce9513612901/go.mod h1:tIUCLVX6utdBPiP6dAjY7GUfwIQngeSMK5A9aRqe6gA= +github.com/LeanerCloud/cloud-commitments-go/pkg v0.0.0-20260929105827-b3b4cb5e3d80 h1:wVKlMokfaME/Lw525Qz3F155R3nY4VmuaIyh+6d4sCU= +github.com/LeanerCloud/cloud-commitments-go/pkg v0.0.0-20260929105827-b3b4cb5e3d80/go.mod h1:ApWBliDXe099f3oDXBz41K/I9v4bHvn1dG/BGoRmHlw= +github.com/LeanerCloud/cloud-commitments-go/providers/aws v0.0.0-20261002152209-006ef5c8d0a2 h1:rYxXq0G0O2lkf8AB1ni9m4pVhIdb6Z0y52sTxS1RKIA= +github.com/LeanerCloud/cloud-commitments-go/providers/aws v0.0.0-20261002152209-006ef5c8d0a2/go.mod h1:d4nsy61/Ptib0SxqMg0ble3yksVm3+QbtUtmGja82PA= github.com/LeanerCloud/cloud-commitments-go/providers/azure v0.0.0-20260928214714-ce9513612901 h1:iSdHYdmGUjcjtSmgGZxptBzDuLrJ9KR2mh1/x7FNvSY= github.com/LeanerCloud/cloud-commitments-go/providers/azure v0.0.0-20260928214714-ce9513612901/go.mod h1:zgCL/ozOkcZUDbEC7a2UwW+6MPcLoBlEU2WUZ0TORUs= github.com/LeanerCloud/cloud-commitments-go/providers/gcp v0.0.0-20260928214714-ce9513612901 h1:i6OwXLUheudN3GfwnYXdKuEq8vPdE9qkNJ+r71LRLKA= diff --git a/internal/scheduler/partial_sweep_eviction_test.go b/internal/scheduler/partial_sweep_eviction_test.go index 7dd0d6e8..dd1e69cd 100644 --- a/internal/scheduler/partial_sweep_eviction_test.go +++ b/internal/scheduler/partial_sweep_eviction_test.go @@ -3,9 +3,11 @@ package scheduler import ( "context" "errors" + "fmt" "testing" "github.com/LeanerCloud/cloud-commitments-go/pkg/common" + awsrecommendations "github.com/LeanerCloud/cloud-commitments-go/providers/aws/recommendations" azureprovider "github.com/LeanerCloud/cloud-commitments-go/providers/azure" "github.com/LeanerCloud/cloud-commitments-platform/internal/config" "github.com/stretchr/testify/assert" @@ -13,6 +15,75 @@ import ( "github.com/stretchr/testify/require" ) +func TestAWSIncompleteSweep(t *testing.T) { + incomplete := &awsrecommendations.IncompleteRecommendationsError{ + FailedDetails: 2, FailedScopes: 1, Causes: []error{errors.New("malformed quantity"), errors.New("API scope failed")}, + } + for _, tc := range []struct { + name string + err error + complete bool + fatal bool + }{ + {name: "complete", complete: true}, + {name: "incomplete", err: incomplete}, + {name: "wrapped incomplete", err: fmt.Errorf("collect: %w", incomplete)}, + {name: "ordinary failure", err: errors.New("API unavailable"), fatal: true}, + {name: "canceled", err: context.Canceled, fatal: true}, + {name: "deadline", err: context.DeadlineExceeded, fatal: true}, + } { + t.Run(tc.name, func(t *testing.T) { + complete, err := tolerateIncompleteSweep("aws", tc.err) + require.Equal(t, tc.complete, complete) + if tc.fatal { + require.ErrorIs(t, err, tc.err) + } else { + require.NoError(t, err) + } + }) + } +} + +func TestAWSIncompleteFallbackCannotAuthorizeEviction(t *testing.T) { + incomplete := &awsrecommendations.IncompleteRecommendationsError{FailedDetails: 1, Causes: []error{errors.New("invalid quantity")}} + for _, tc := range []struct { + name string + firstErr error + lastErr error + }{ + {name: "all-invalid then complete", firstErr: incomplete}, + {name: "empty then incomplete", lastErr: incomplete}, + {name: "both incomplete", firstErr: incomplete, lastErr: incomplete}, + } { + t.Run(tc.name, func(t *testing.T) { + client := new(MockRecommendationsClient) + client.On("GetAllRecommendations", mock.Anything).Return([]common.Recommendation{}, tc.firstErr).Once() + for _, service := range []common.ServiceType{common.ServiceEC2, common.ServiceRDS, common.ServiceElastiCache, common.ServiceOpenSearch, common.ServiceRedshift, common.ServiceSavingsPlansAll} { + var rows []common.Recommendation + var lastErr error + if service == common.ServiceRDS { + rows = []common.Recommendation{{Service: common.ServiceRDS, Region: "us-east-1", ResourceType: "db.t3.medium", Count: 2, + Term: "1yr", PaymentOption: "no-upfront"}} + lastErr = tc.lastErr + } + client.On("GetRecommendations", mock.Anything, mock.MatchedBy(func(params *common.RecommendationParams) bool { + return params.Service == service && params.Term == "1yr" && params.PaymentOption == "no-upfront" && params.LookbackPeriod == "30d" + })).Return(rows, lastErr).Once() + } + prov := new(MockProvider) + prov.On("GetRecommendationsClient", mock.Anything).Return(client, nil).Once() + s := &Scheduler{} + recs, complete, err := s.fetchAndConvert(context.Background(), prov, "aws", nil, + &config.GlobalConfig{DefaultTerm: 1, DefaultPayment: "no-upfront", RecommendationsLookbackDays: 30}) + require.NoError(t, err) + require.Len(t, recs, 1) + require.False(t, complete) + client.AssertExpectations(t) + prov.AssertExpectations(t) + }) + } +} + // A partially-swept account must never authorize stale-row eviction. // // UpsertRecommendations deletes an account's previous-cycle rows with diff --git a/internal/scheduler/recommendation_completeness_integration_test.go b/internal/scheduler/recommendation_completeness_integration_test.go new file mode 100644 index 00000000..a679635c --- /dev/null +++ b/internal/scheduler/recommendation_completeness_integration_test.go @@ -0,0 +1,246 @@ +//go:build integration + +package scheduler + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net/http" + "path/filepath" + "strings" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/LeanerCloud/cloud-commitments-go/pkg/provider" + awsprovider "github.com/LeanerCloud/cloud-commitments-go/providers/aws" + "github.com/LeanerCloud/cloud-commitments-platform/internal/config" + "github.com/LeanerCloud/cloud-commitments-platform/internal/database/postgres/migrations" + "github.com/LeanerCloud/cloud-commitments-platform/internal/database/postgres/testhelpers" + "github.com/aws/aws-sdk-go-v2/aws" + "github.com/google/uuid" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" +) + +type completenessHTTP struct { + resourceType string + details []string + fallback []string + failScope bool + rdsCalls atomic.Int32 + unexpected atomic.Int32 + mu sync.Mutex + requests []completenessRequest +} + +type completenessRequest struct { + Operation string + Service string + LookbackPeriodInDays string + TermInYears string + PaymentOption string + SavingsPlansType string +} + +func (f *completenessHTTP) Do(req *http.Request) (*http.Response, error) { + var input completenessRequest + if err := json.NewDecoder(req.Body).Decode(&input); err != nil { + return nil, err + } + input.Operation = req.Header.Get("X-Amz-Target") + f.mu.Lock() + f.requests = append(f.requests, input) + f.mu.Unlock() + body := any(map[string]any{}) + status := http.StatusOK + switch req.Header.Get("X-Amz-Target") { + case "AWSInsightsIndexService.GetReservationPurchaseRecommendation": + allowed := map[string]bool{ + "Amazon Elastic Compute Cloud - Compute": true, "Amazon Relational Database Service": true, + "Amazon ElastiCache": true, "Amazon OpenSearch Service": true, "Amazon Redshift": true, + } + if !allowed[input.Service] { + f.unexpected.Add(1) + return nil, fmt.Errorf("invalid Cost Explorer Service %q", input.Service) + } + details := []any{} + if input.Service == "Amazon Relational Database Service" { + call := f.rdsCalls.Add(1) + quantities := f.details + if input.LookbackPeriodInDays == "THIRTY_DAYS" { + quantities = f.fallback + } + for _, quantity := range quantities { + details = append(details, map[string]any{ + "RecommendedNumberOfInstancesToPurchase": quantity, + "EstimatedMonthlySavingsAmount": "10", "EstimatedMonthlyOnDemandCost": "30", + "InstanceDetails": map[string]any{"RDSInstanceDetails": map[string]any{ + "InstanceType": f.resourceType, "Region": "us-east-1", "DeploymentOption": "Single-AZ", + }}, + }) + } + if f.failScope && call == 1 { + status = http.StatusBadRequest + } + } + body = map[string]any{"Recommendations": []any{map[string]any{"RecommendationDetails": details}}} + case "AWSInsightsIndexService.GetSavingsPlansPurchaseRecommendation", "AWSInsightsIndexService.GetReservationCoverage": + default: + f.unexpected.Add(1) + return nil, fmt.Errorf("unexpected AWS operation %s", req.Header.Get("X-Amz-Target")) + } + if status != http.StatusOK { + body = map[string]any{"__type": "AccessDeniedException", "message": "synthetic scope failure"} + } + raw, err := json.Marshal(body) + if err != nil { + return nil, err + } + return &http.Response{StatusCode: status, Header: http.Header{"Content-Type": {"application/x-amz-json-1.1"}}, + Body: io.NopCloser(strings.NewReader(string(raw))), Request: req}, nil +} + +func completenessProvider(t *testing.T, fixture *completenessHTTP) *MockProvider { + t.Helper() + client := awsprovider.NewRecommendationsClient(aws.Config{ + Region: "us-east-1", HTTPClient: fixture, Retryer: func() aws.Retryer { return aws.NopRetryer{} }, + Credentials: aws.CredentialsProviderFunc(func(context.Context) (aws.Credentials, error) { + return aws.Credentials{AccessKeyID: "synthetic", SecretAccessKey: "synthetic"}, nil + }), + }) + prov := new(MockProvider) + prov.On("GetRecommendationsClient", mock.Anything).Return(client, nil).Once() + t.Cleanup(func() { + prov.AssertExpectations(t) + require.Zero(t, fixture.unexpected.Load()) + require.Positive(t, fixture.rdsCalls.Load()) + }) + return prov +} + +func TestAWSRecommendationCompletenessPersistence(t *testing.T) { + for _, tc := range []struct { + name string + registered bool + details []string + fallback []string + failScope bool + complete bool + wantRows int + }{ + {name: "ambient mixed", details: []string{"2", "invalid"}, wantRows: 6}, + {name: "ambient all invalid", details: []string{"invalid"}}, + {name: "ambient valid", details: []string{"2"}, complete: true, wantRows: 6}, + {name: "ambient empty", complete: true}, + {name: "configured lookback returns valid fallback", fallback: []string{"2"}, complete: true, wantRows: 1}, + {name: "configured lookback remains empty", fallback: []string{}, complete: true}, + {name: "ambient all invalid then clean fallback", details: []string{"invalid"}, fallback: []string{"2"}, wantRows: 1}, + {name: "ambient empty then invalid fallback", fallback: []string{"invalid"}}, + {name: "ambient failed API scope", details: []string{"2"}, failScope: true, wantRows: 5}, + {name: "registered mixed", registered: true, details: []string{"2", "invalid"}, wantRows: 6}, + {name: "registered all invalid", registered: true, details: []string{"invalid"}}, + } { + t.Run(tc.name, func(t *testing.T) { + ctx := context.Background() + container, err := testhelpers.SetupPostgresContainer(ctx, t) + require.NoError(t, err) + t.Cleanup(func() { container.Cleanup(ctx) }) + migrationPath, err := filepath.Abs("../database/postgres/migrations") + require.NoError(t, err) + require.NoError(t, migrations.RunMigrations(ctx, container.DB.Pool(), migrationPath, "", "")) + store := config.NewPostgresStore(container.DB) + factory := new(MockProviderFactory) + fixture := &completenessHTTP{resourceType: "db.t3.medium", details: tc.details, fallback: tc.fallback, failScope: tc.failScope} + factory.On("CreateAndValidateProvider", mock.Anything, "aws", (*provider.ProviderConfig)(nil)). + Return(completenessProvider(t, fixture), nil).Once() + var accountID *string + var completeAccountID string + if tc.registered { + id := uuid.NewString() + accountID = &id + completeAccountID = uuid.NewString() + for _, acct := range []config.CloudAccount{ + {ID: id, Name: "incomplete", Provider: "aws", ExternalID: "111111111111", Enabled: true, AWSAuthMode: "role_arn"}, + {ID: completeAccountID, Name: "complete", Provider: "aws", ExternalID: "222222222222", Enabled: true, AWSAuthMode: "role_arn", AWSRoleARN: "arn:aws:iam::222222222222:role/synthetic"}, + } { + require.NoError(t, store.CreateCloudAccount(ctx, &acct)) + } + factory.On("CreateAndValidateProvider", mock.Anything, "aws", mock.MatchedBy(func(cfg *provider.ProviderConfig) bool { return cfg != nil })). + Return(completenessProvider(t, &completenessHTTP{resourceType: "db.r5.large", details: []string{"2"}}), nil).Once() + } + seed := []config.RecommendationRecord{ + {ID: "missing-offer", Provider: "aws", CloudAccountID: accountID, Service: "rds", Region: "us-east-1", ResourceType: "db.t3.large", Savings: 40, Count: 1, Term: 12, Payment: "no-upfront"}, + {ID: "unswept-provider", Provider: "azure", Service: "vm", Region: "eastus", ResourceType: "D2", Savings: 20, Count: 1, Term: 12, Payment: "no-upfront"}, + } + if tc.registered { + seed = append(seed, config.RecommendationRecord{ID: "complete-stale", Provider: "aws", CloudAccountID: &completeAccountID, + Service: "rds", Region: "us-east-1", ResourceType: "db.t3.large", Savings: 40, Count: 1, Term: 12, Payment: "no-upfront"}) + } + require.NoError(t, store.UpsertRecommendations(ctx, time.Now().Add(-time.Hour), seed, nil)) + s := &Scheduler{config: store, providerFactory: factory} + var globalCfg *config.GlobalConfig + if tc.fallback != nil { + globalCfg = &config.GlobalConfig{DefaultTerm: 1, DefaultPayment: "no-upfront", RecommendationsLookbackDays: 30} + } + recs, ids, err := s.collectAWSRecommendations(ctx, globalCfg) + require.NoError(t, err) + wantRows := tc.wantRows + switch { + case tc.registered: + wantRows += 6 + require.Equal(t, []string{completeAccountID}, ids) + case tc.complete: + require.Equal(t, []string{""}, ids) + default: + require.Empty(t, ids) + } + require.Len(t, recs, wantRows) + s.persistCollection(ctx, recs, expandSuccessfulCollects("aws", ids), nil) + rows, err := store.ListStoredRecommendations(ctx, config.RecommendationFilter{}) + require.NoError(t, err) + byID := make(map[string]config.RecommendationRecord, len(rows)) + for _, row := range rows { + byID[row.ID] = row + } + _, missing := byID["missing-offer"] + require.Equal(t, !tc.complete, missing, "incomplete collection must retain previous offers") + require.Contains(t, byID, "unswept-provider") + require.NotContains(t, byID, "complete-stale") + for _, rec := range recs { + require.Contains(t, byID, rec.ID, "surviving SDK recommendations must be persisted") + require.Equal(t, rec.CloudAccountID, byID[rec.ID].CloudAccountID) + } + fallbackServices := []string{} + for _, request := range fixture.requests { + if request.Operation == "AWSInsightsIndexService.GetReservationCoverage" { + continue + } + require.Contains(t, []string{"SEVEN_DAYS", "THIRTY_DAYS"}, request.LookbackPeriodInDays) + if request.LookbackPeriodInDays == "SEVEN_DAYS" { + continue + } + require.Equal(t, "ONE_YEAR", request.TermInYears) + require.Equal(t, "NO_UPFRONT", request.PaymentOption) + service := request.Service + if request.Operation == "AWSInsightsIndexService.GetSavingsPlansPurchaseRecommendation" { + service = "SP:" + request.SavingsPlansType + } + fallbackServices = append(fallbackServices, service) + } + if tc.fallback != nil { + require.Equal(t, int32(7), fixture.rdsCalls.Load()) + require.ElementsMatch(t, []string{ + "Amazon Elastic Compute Cloud - Compute", "Amazon Relational Database Service", "Amazon ElastiCache", + "Amazon OpenSearch Service", "Amazon Redshift", "SP:COMPUTE_SP", "SP:EC2_INSTANCE_SP", "SP:SAGEMAKER_SP", "SP:DATABASE_SP", + }, fallbackServices) + } else { + require.Empty(t, fallbackServices) + } + factory.AssertExpectations(t) + }) + } +} diff --git a/internal/scheduler/scheduler.go b/internal/scheduler/scheduler.go index c9ba275e..9b831683 100644 --- a/internal/scheduler/scheduler.go +++ b/internal/scheduler/scheduler.go @@ -3,6 +3,7 @@ package scheduler import ( "context" + "errors" "fmt" "os" "sort" @@ -16,6 +17,7 @@ import ( "github.com/LeanerCloud/cloud-commitments-go/pkg/concurrency" "github.com/LeanerCloud/cloud-commitments-go/pkg/logging" "github.com/LeanerCloud/cloud-commitments-go/pkg/provider" + awsrecommendations "github.com/LeanerCloud/cloud-commitments-go/providers/aws/recommendations" azureprovider "github.com/LeanerCloud/cloud-commitments-go/providers/azure" gcpprovider "github.com/LeanerCloud/cloud-commitments-go/providers/gcp" "github.com/LeanerCloud/cloud-commitments-platform/internal/config" @@ -939,7 +941,7 @@ func (s *Scheduler) enabledAccounts(ctx context.Context, providerName string) [] return accounts } -// tolerateIncompleteSweep converts an org-wide partial-subscription failure +// tolerateIncompleteSweep converts an incomplete recommendation collection // into a warning and a nil error, so the caller keeps the recommendations that // WERE collected instead of discarding them. // @@ -984,6 +986,12 @@ func (s *Scheduler) enabledAccounts(ctx context.Context, providerName string) [] // design decision and is left as follow-up work. With eviction withheld the // residue is observability, not lost rows. func tolerateIncompleteSweep(providerName string, err error) (complete bool, _ error) { + var incomplete *awsrecommendations.IncompleteRecommendationsError + if errors.As(err, &incomplete) { + logging.Warnf("%s recommendations incomplete: failed_details=%d failed_scopes=%d; %v; keeping collected rows and withholding stale-row eviction", + providerName, incomplete.FailedDetails, incomplete.FailedScopes, incomplete) + return false, nil + } partial := azureprovider.AsPartialSubscriptionFailure(err) if partial == nil { return err == nil, err @@ -1016,26 +1024,22 @@ func (s *Scheduler) fetchAndConvert(ctx context.Context, prov provider.Provider, if lookbackDays == 0 { lookbackDays = config.DefaultRecommendationsLookbackDays } - params := common.RecommendationParams{ - Term: fmt.Sprintf("%dyr", globalCfg.DefaultTerm), - PaymentOption: globalCfg.DefaultPayment, - LookbackPeriod: fmt.Sprintf("%dd", lookbackDays), - } - var recErr error - recs, recErr = recClient.GetRecommendations(ctx, ¶ms) - fallbackComplete, recErr := tolerateIncompleteSweep(providerName, recErr) - if recErr != nil { - // Fail loud: a misconfigured DefaultPayment/DefaultTerm or a CE - // failure on this fallback must surface to the operator instead - // of silently presenting as "zero recommendations". - return nil, false, fmt.Errorf("failed to get %s recommendations with default term/payment fallback (term=%s, payment=%s, lookback=%s): %w", - providerName, params.Term, params.PaymentOption, params.LookbackPeriod, recErr) + for _, service := range []common.ServiceType{common.ServiceEC2, common.ServiceRDS, common.ServiceElastiCache, common.ServiceOpenSearch, common.ServiceRedshift, common.ServiceSavingsPlansAll} { + params := common.RecommendationParams{ + Service: service, + Term: fmt.Sprintf("%dyr", globalCfg.DefaultTerm), + PaymentOption: globalCfg.DefaultPayment, + LookbackPeriod: fmt.Sprintf("%dd", lookbackDays), + } + fallbackRecs, recErr := recClient.GetRecommendations(ctx, ¶ms) + fallbackComplete, recErr := tolerateIncompleteSweep(providerName, recErr) + if recErr != nil { + return nil, false, fmt.Errorf("failed to get %s recommendations with default term/payment fallback (service=%s, term=%s, payment=%s, lookback=%s): %w", + providerName, params.Service, params.Term, params.PaymentOption, params.LookbackPeriod, recErr) + } + recs = append(recs, fallbackRecs...) + complete = complete && fallbackComplete } - // AND, not assignment: the fallback re-queries the same scope, so an - // incomplete first sweep stays incomplete even if the retry happens to - // come back whole. Overwriting here would let the retry launder away - // the first sweep's missing subscriptions and re-authorize eviction. - complete = complete && fallbackComplete } result := s.convertRecommendations(recs, providerName) if accountID != nil { diff --git a/internal/scheduler/scheduler_test.go b/internal/scheduler/scheduler_test.go index 7c450c49..9771ea21 100644 --- a/internal/scheduler/scheduler_test.go +++ b/internal/scheduler/scheduler_test.go @@ -2108,7 +2108,15 @@ func TestScheduler_CollectAWSRecommendations_FallbackToFiltered(t *testing.T) { mockFactory.On("CreateAndValidateProvider", mock.Anything, "aws", mock.Anything).Return(mockProvider, nil) mockProvider.On("GetRecommendationsClient", ctx).Return(mockRecClient, nil) mockRecClient.On("GetAllRecommendations", ctx).Return([]common.Recommendation{}, nil) // Empty - mockRecClient.On("GetRecommendations", ctx, mock.AnythingOfType("*common.RecommendationParams")).Return(filteredRecommendations, nil) + for _, service := range []common.ServiceType{common.ServiceEC2, common.ServiceRDS, common.ServiceElastiCache, common.ServiceOpenSearch, common.ServiceRedshift, common.ServiceSavingsPlansAll} { + var rows []common.Recommendation + if service == common.ServiceEC2 { + rows = filteredRecommendations + } + mockRecClient.On("GetRecommendations", ctx, mock.MatchedBy(func(params *common.RecommendationParams) bool { + return params.Service == service && params.Term == "3yr" && params.PaymentOption == "all-upfront" && params.LookbackPeriod == "7d" + })).Return(rows, nil).Once() + } scheduler := &Scheduler{ config: mockStore, @@ -2118,6 +2126,7 @@ func TestScheduler_CollectAWSRecommendations_FallbackToFiltered(t *testing.T) { recs, _, err := scheduler.collectAWSRecommendations(ctx, globalCfg) require.NoError(t, err) assert.Len(t, recs, 1) + mockRecClient.AssertExpectations(t) } // Regression test for COR-05 (#1168): when the primary sweep returns zero @@ -2146,7 +2155,9 @@ func TestScheduler_CollectAWSRecommendations_FallbackError(t *testing.T) { mockFactory.On("CreateAndValidateProvider", mock.Anything, "aws", mock.Anything).Return(mockProvider, nil) mockProvider.On("GetRecommendationsClient", ctx).Return(mockRecClient, nil) mockRecClient.On("GetAllRecommendations", ctx).Return([]common.Recommendation{}, nil) // Empty -> triggers fallback - mockRecClient.On("GetRecommendations", ctx, mock.AnythingOfType("*common.RecommendationParams")).Return(nil, fallbackErr) + mockRecClient.On("GetRecommendations", ctx, mock.MatchedBy(func(params *common.RecommendationParams) bool { + return params.Service == common.ServiceEC2 && params.Term == "3yr" && params.PaymentOption == "all-upfront" && params.LookbackPeriod == "7d" + })).Return(nil, fallbackErr).Once() scheduler := &Scheduler{ config: mockStore,