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
4 changes: 2 additions & 2 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 4 additions & 4 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
71 changes: 71 additions & 0 deletions internal/scheduler/partial_sweep_eviction_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,16 +3,87 @@ 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"
"github.com/stretchr/testify/mock"
"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
Expand Down
246 changes: 246 additions & 0 deletions internal/scheduler/recommendation_completeness_integration_test.go
Original file line number Diff line number Diff line change
@@ -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)
})
}
}
Loading
Loading