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
3 changes: 3 additions & 0 deletions providers/aws/recommendations/coverage.go
Original file line number Diff line number Diff line change
Expand Up @@ -263,6 +263,9 @@ func (c *Client) fetchCoveragePaged(
) error {
var token *string
for {
if err := ctx.Err(); err != nil {
return fmt.Errorf("coverage: pagination cancelled: %w", err)
}
input.NextPageToken = token
result, err := c.fetchCoveragePage(ctx, input)
if err != nil {
Expand Down
27 changes: 27 additions & 0 deletions providers/aws/recommendations/coverage_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -361,6 +361,33 @@ func TestNormaliseRDSEngine(t *testing.T) {
}
}

// TestFetchCoveragePaged_CtxCancelReturnsError asserts that a cancelled context
// is treated as a hard stop inside the pagination loop and surfaces an error
// rather than returning a partial/empty result silently
// (feedback_ctx_cancel_terminal).
func TestFetchCoveragePaged_CtxCancelReturnsError(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel() // cancel before the call

// The mock returns a non-nil token on the first call so the loop
// would try to paginate -- but ctx.Err() must fire first.
token := "page2"
mock := &mockCoverageCE{
coverageOutput: &costexplorer.GetReservationCoverageOutput{
NextPageToken: &token,
},
}
client := NewClientWithAPI(mock, "us-east-1")

err := client.fetchCoveragePaged(
ctx,
&costexplorer.GetReservationCoverageInput{},
func(_, _ string, _ PoolCoverage) {},
720,
)
require.Error(t, err, "cancelled context must surface an error from fetchCoveragePaged")
}

// TestNormaliseDeployment locks deployment-option canonicalisation. CE
// returns "Single-AZ" / "Multi-AZ" / "Multi-AZ (readable standbys)" while
// the parser stores "single-az" / "multi-az"; both forms must collapse to
Expand Down
10 changes: 8 additions & 2 deletions providers/aws/recommendations/utilization.go
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,9 @@ func (c *Client) GetRIUtilization(ctx context.Context, lookbackDays int) ([]RIUt

var nextPageToken *string
for {
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("utilization: pagination cancelled: %w", err)
}
input.NextPageToken = nextPageToken

result, err := c.fetchUtilizationPage(ctx, input)
Expand All @@ -72,6 +75,10 @@ func (c *Client) GetRIUtilization(ctx context.Context, lookbackDays int) ([]RIUt
nextPageToken = result.NextPageToken
}

return buildUtilizations(agg), nil
}

func buildUtilizations(agg map[string]*riAccumulator) []RIUtilization {
utilizations := make([]RIUtilization, 0, len(agg))
for id, a := range agg {
pct := 0.0
Expand All @@ -86,8 +93,7 @@ func (c *Client) GetRIUtilization(ctx context.Context, lookbackDays int) ([]RIUt
UnusedHours: a.unusedHours,
})
}

return utilizations, nil
return utilizations
}

// fetchUtilizationPage calls the Cost Explorer API with rate-limit retry.
Expand Down
47 changes: 40 additions & 7 deletions providers/azure/services/cache/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,15 @@ import (
"github.com/LeanerCloud/CUDly/providers/azure/services/internal/reservations"
)

// maxRecsPages caps Consumption API recommendation pagination.
const maxRecsPages = 10

// maxReservationsPages caps reservation-detail pagination.
const maxReservationsPages = 50

// maxCachesPages caps Redis cache list pagination.
const maxCachesPages = 20

// redisSKUEntry holds the SKU-catalogue-derived fields the converter
// wants for each Redis SKU. Sourced from the cache's Properties:
// - shardCount: Properties.ShardCount (Premium-tier clustered caches).
Expand Down Expand Up @@ -166,7 +175,13 @@ func (c *CacheClient) GetRecommendations(ctx context.Context, params common.Reco
pager = client.NewListPager(scope, &armconsumption.ReservationRecommendationsClientListOptions{Filter: &filter})
}

for pager.More() {
for pageIdx := 0; pager.More(); pageIdx++ {
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("context cancelled during pagination: %w", err)
}
if pageIdx >= maxRecsPages {
return nil, fmt.Errorf("cache: GetRecommendations pagination cap (%d pages) reached", maxRecsPages)
}
page, err := pager.NextPage(ctx)
if err != nil {
return nil, fmt.Errorf("failed to get Redis Cache recommendations: %w", err)
Expand Down Expand Up @@ -216,7 +231,13 @@ func (c *CacheClient) createReservationsPager() (ReservationsDetailsPager, error
func (c *CacheClient) collectRedisReservations(ctx context.Context, pager ReservationsDetailsPager) ([]common.Commitment, error) {
commitments := make([]common.Commitment, 0)

for pager.More() {
for pageIdx := 0; pager.More(); pageIdx++ {
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("context cancelled during pagination: %w", err)
}
if pageIdx >= maxReservationsPages {
return nil, fmt.Errorf("cache: GetExistingCommitments pagination cap (%d pages) reached", maxReservationsPages)
}
page, err := pager.NextPage(ctx)
if err != nil {
return nil, fmt.Errorf("cache: list reservations: %w", err)
Expand Down Expand Up @@ -403,7 +424,10 @@ func (c *CacheClient) GetValidResourceTypes(ctx context.Context) ([]string, erro
return c.getCommonSKUs(), nil
}

skuSet := c.collectSKUsFromCaches(ctx, pager)
skuSet, err := c.collectSKUsFromCaches(ctx, pager)
if err != nil {
return nil, err
}

// If we found SKUs from existing caches, use those
if len(skuSet) > 0 {
Expand All @@ -429,11 +453,20 @@ func (c *CacheClient) createRedisCachesPager() (RedisCachesPager, error) {
return client.NewListBySubscriptionPager(nil), nil
}

// collectSKUsFromCaches collects SKUs from existing Redis caches
func (c *CacheClient) collectSKUsFromCaches(ctx context.Context, pager RedisCachesPager) map[string]bool {
// collectSKUsFromCaches collects SKUs from existing Redis caches.
// Returns (nil, err) on context cancellation so callers can propagate the error
// instead of silently using a partial result set.
func (c *CacheClient) collectSKUsFromCaches(ctx context.Context, pager RedisCachesPager) (map[string]bool, error) {
skuSet := make(map[string]bool)

for pager.More() {
for pageIdx := 0; pager.More(); pageIdx++ {
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("cache: GetValidResourceTypes context cancelled after %d pages: %w", pageIdx, err)
}
if pageIdx >= maxCachesPages {
log.Printf("WARNING: cache: GetValidResourceTypes pagination cap (%d pages) reached", maxCachesPages)
break
}
page, err := pager.NextPage(ctx)
if err != nil {
// If we can't list existing caches, fall back to known SKU families
Expand All @@ -447,7 +480,7 @@ func (c *CacheClient) collectSKUsFromCaches(ctx context.Context, pager RedisCach
}
}

return skuSet
return skuSet, nil
}

// extractSKUFromCache extracts the full SKU name from a cache resource
Expand Down
42 changes: 40 additions & 2 deletions providers/azure/services/cache/client_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -188,7 +188,7 @@ func TestCacheClient_GetValidResourceTypes_Fallback(t *testing.T) {
// When API calls fail, GetValidResourceTypes should return common SKUs
client := NewClient(nil, "invalid-subscription", "eastus")

skus, err := client.GetValidResourceTypes(nil)
skus, err := client.GetValidResourceTypes(context.Background())
require.NoError(t, err)
require.NotEmpty(t, skus)

Expand All @@ -204,7 +204,7 @@ func TestCacheClient_ValidateOffering_InvalidSKU(t *testing.T) {
ResourceType: "InvalidSKU_X99",
}

err := client.ValidateOffering(nil, rec)
err := client.ValidateOffering(context.Background(), rec)
assert.Error(t, err)
assert.Contains(t, err.Error(), "invalid Azure Redis Cache SKU")
}
Expand Down Expand Up @@ -1196,3 +1196,41 @@ func TestCacheClient_PurchaseCommitment_DisplayNameConformsToAzureAllowlist(t *t
assert.Regexp(t, `^redis-`, capturedDisplayName)
assert.Contains(t, capturedDisplayName, "Premium_P1")
}

// infiniteRedisCachesPager is a pager that always reports More()=true,
// used to exercise the maxCachesPages budget cap.
type infiniteRedisCachesPager struct{}

func (p *infiniteRedisCachesPager) More() bool { return true }
func (p *infiniteRedisCachesPager) NextPage(_ context.Context) (armredis.ClientListBySubscriptionResponse, error) {
return armredis.ClientListBySubscriptionResponse{}, nil
}

// TestCacheClient_GetValidResourceTypes_CtxCancelReturnsError asserts that a
// cancelled context is treated as a hard stop and surfaces an error rather than
// returning a silent partial result (feedback_ctx_cancel_terminal).
func TestCacheClient_GetValidResourceTypes_CtxCancelReturnsError(t *testing.T) {
ctx, cancel := context.WithCancel(context.Background())
cancel() // cancel immediately

client := NewClient(nil, "test-subscription", "eastus")
client.SetRedisCachesPager(&infiniteRedisCachesPager{})

_, err := client.GetValidResourceTypes(ctx)
require.Error(t, err, "cancelled context must produce an error, not a silent partial result")
}

// TestCacheClient_GetValidResourceTypes_PageCapReturnsError asserts that the
// page budget is enforced: an unbounded pager is stopped after maxCachesPages
// iterations and the function returns an error rather than looping forever.
func TestCacheClient_GetValidResourceTypes_PageCapFires(t *testing.T) {
client := NewClient(nil, "test-subscription", "eastus")
client.SetRedisCachesPager(&infiniteRedisCachesPager{})

// GetValidResourceTypes falls back to common SKUs after hitting the log-only
// cap, so the call itself succeeds. The important invariant is that it
// terminates rather than looping forever.
skus, err := client.GetValidResourceTypes(context.Background())
require.NoError(t, err)
require.NotEmpty(t, skus, "page cap must trigger fallback to common SKUs, not an empty result")
}
49 changes: 44 additions & 5 deletions providers/azure/services/compute/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -58,6 +58,19 @@ type HTTPClient interface {
// Either field stays at the zero value when the capability is missing
// or unparseable; common.ComputeDetails treats 0 as "unknown" (the
// JSON tags on VCPU/MemoryGB are omitempty).

// maxRecsPages caps Consumption API recommendation pagination to avoid
// burning a Lambda deadline on a stalled or unexpectedly deep result set.
const maxRecsPages = 10

// maxReservationsPages caps reservation-detail pagination.
// Large orgs may have hundreds of reservations spread over many pages.
const maxReservationsPages = 50

// maxSKUPages caps Azure ResourceSKUs pagination.
// The SKU catalogue for a subscription can run to many pages.
const maxSKUPages = 20

type vmSKUEntry struct {
vCPUs int
memoryGB float64
Expand Down Expand Up @@ -185,7 +198,13 @@ func (c *ComputeClient) GetRecommendations(ctx context.Context, params common.Re
pager = client.NewListPager(scope, &armconsumption.ReservationRecommendationsClientListOptions{Filter: &filter})
}

for pager.More() {
for pageIdx := 0; pager.More(); pageIdx++ {
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("context cancelled during pagination: %w", err)
}
if pageIdx >= maxRecsPages {
return nil, fmt.Errorf("compute: GetRecommendations pagination cap (%d pages) reached", maxRecsPages)
}
page, err := pager.NextPage(ctx)
if err != nil {
return nil, fmt.Errorf("failed to get VM recommendations: %w", err)
Expand Down Expand Up @@ -238,7 +257,13 @@ func (c *ComputeClient) createReservationsPager() (ReservationsDetailsPager, err
func (c *ComputeClient) collectVMReservations(ctx context.Context, pager ReservationsDetailsPager) ([]common.Commitment, error) {
commitments := make([]common.Commitment, 0)

for pager.More() {
for pageIdx := 0; pager.More(); pageIdx++ {
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("context cancelled during pagination: %w", err)
}
if pageIdx >= maxReservationsPages {
return nil, fmt.Errorf("compute: GetExistingCommitments pagination cap (%d pages) reached", maxReservationsPages)
}
page, err := pager.NextPage(ctx)
if err != nil {
return nil, fmt.Errorf("compute: list reservations: %w", err)
Expand Down Expand Up @@ -553,7 +578,13 @@ func (c *ComputeClient) createResourceSKUsPager() (ResourceSKUsPager, error) {
func (c *ComputeClient) collectVMSizesFromSKUs(ctx context.Context, pager ResourceSKUsPager) ([]string, error) {
vmSizes := make([]string, 0)

for pager.More() {
for pageIdx := 0; pager.More(); pageIdx++ {
if err := ctx.Err(); err != nil {
return nil, fmt.Errorf("context cancelled during pagination: %w", err)
}
if pageIdx >= maxSKUPages {
return nil, fmt.Errorf("compute: GetValidResourceTypes pagination cap (%d pages) reached", maxSKUPages)
}
page, err := pager.NextPage(ctx)
if err != nil {
return nil, fmt.Errorf("failed to list VM sizes: %w", err)
Expand Down Expand Up @@ -777,10 +808,18 @@ func (c *ComputeClient) fetchSKUCatalogue(ctx context.Context) map[string]vmSKUE
return nil
}
out := make(map[string]vmSKUEntry)
for pager.More() {
for pageIdx := 0; pager.More(); pageIdx++ {
if err := ctx.Err(); err != nil {
logging.Warnf("azure compute: SKU catalogue fetch cancelled for region %s after %d pages: %v; partial cache discarded", c.region, pageIdx, err)
return nil
}
if pageIdx >= maxSKUPages {
logging.Warnf("azure compute: SKU catalogue pagination cap (%d pages) reached for region %s; partial cache (%d entries) used", maxSKUPages, c.region, len(out))
break
}
page, err := pager.NextPage(ctx)
if err != nil {
logging.Warnf("azure compute: SKU catalogue page fetch failed for region %s: %v — partial cache (%d entries) discarded, Details.VCPU/MemoryGB left at 0", c.region, err, len(out))
logging.Warnf("azure compute: SKU catalogue page fetch failed for region %s: %v; partial cache (%d entries) discarded, Details.VCPU/MemoryGB left at 0", c.region, err, len(out))
return nil
}
c.populateVMSKUMapFromPage(out, page.Value)
Expand Down
Loading
Loading