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
21 changes: 20 additions & 1 deletion cmd/multi_service_helpers.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ import (
"github.com/LeanerCloud/CUDly/pkg/common"
"github.com/LeanerCloud/CUDly/pkg/provider"
"github.com/LeanerCloud/CUDly/providers/aws/recommendations"
azureprovider "github.com/LeanerCloud/CUDly/providers/azure"
"github.com/aws/aws-sdk-go-v2/aws"
awsec2 "github.com/aws/aws-sdk-go-v2/service/ec2"
)
Expand Down Expand Up @@ -99,7 +100,13 @@ func getAllAWSRegionsWithClient(ctx context.Context, ec2Client EC2ClientInterfac
// discoverRegionsForService discovers regions that have recommendations for a specific service.
func discoverRegionsForService(ctx context.Context, client provider.RecommendationsClient, service common.ServiceType) ([]string, error) {
recs, err := client.GetRecommendationsForService(ctx, service)
if err != nil {
if partial := azureprovider.AsPartialSubscriptionFailure(err); partial != nil {
// Region discovery is best-effort: the subscriptions that answered
// still tell us where to look. Report the gap rather than dropping
// the discovered regions or failing outright.
AppLogger.Printf(" ⚠️ Region discovery incomplete: %d of %d Azure subscriptions succeeded\n",
partial.Succeeded, partial.Attempted)
} else if err != nil {
return nil, err
}

Expand Down Expand Up @@ -439,6 +446,18 @@ func fetchRecommendationsForRegion(
}

recs, err := recClient.GetRecommendations(ctx, &params)
if partial := azureprovider.AsPartialSubscriptionFailure(err); partial != nil {
// Keep the subscriptions that did answer, but say plainly that the
// sweep was incomplete: without this the operator would read a short
// list as "little to buy here" rather than "some subscriptions were
// never queried".
AppLogger.Printf(" ⚠️ Incomplete: %d of %d Azure subscriptions succeeded; %d could not be queried\n",
partial.Succeeded, partial.Attempted, len(partial.Failed))
for _, f := range partial.Failed {
AppLogger.Printf(" subscription %s: %v\n", f.SubscriptionID, f.Err)
}
return recs
}
if err != nil {
AppLogger.Printf(" ❌ Failed to fetch recommendations: %v\n", err)
return nil
Expand Down
51 changes: 51 additions & 0 deletions internal/scheduler/scheduler.go
Original file line number Diff line number Diff line change
Expand Up @@ -687,6 +687,7 @@ func (s *Scheduler) collectAzureAmbient(ctx context.Context, subscriptionID stri
return nil, fmt.Errorf("get Azure recommendations client: %w", err)
}
recs, err := recClient.GetAllRecommendations(ctx)
err = tolerateIncompleteSweep("azure", err)
if err != nil {
return nil, fmt.Errorf("get Azure recommendations: %w", err)
}
Expand All @@ -706,6 +707,7 @@ func (s *Scheduler) collectGCPAmbient(ctx context.Context) ([]config.Recommendat
return nil, fmt.Errorf("get GCP recommendations client: %w", err)
}
recs, err := recClient.GetAllRecommendations(ctx)
err = tolerateIncompleteSweep("gcp", err)
if err != nil {
return nil, fmt.Errorf("get GCP recommendations: %w", err)
}
Expand Down Expand Up @@ -770,6 +772,17 @@ func (s *Scheduler) collectAzureRecommendations(ctx context.Context, _ *config.G
}

func (s *Scheduler) collectAzureForAccount(ctx context.Context, acct config.CloudAccount) ([]config.RecommendationRecord, error) {
// Every recommendation returned below is tagged with THIS account's UUID
// (see tagAccount at the end of this function), so the provider must be
// pinned to this account's subscription. An empty AzureSubscriptionID
// leaves the provider unpinned, and an unpinned provider now fans out
// across every subscription the credential can see -- which would file
// other subscriptions' recommendations under this account and expose them
// to anyone authorized for it. Fail loud instead; the row is misconfigured.
if acct.AzureSubscriptionID == "" {
return nil, fmt.Errorf("cloud account %s has no azure_subscription_id configured", acct.ID)
}

azCred, err := credentials.ResolveAzureTokenCredentialWithOpts(ctx, &acct, s.credStore, credentials.AzureResolveOptions{
Signer: s.oidcSigner,
IssuerURL: s.oidcIssuerURL,
Expand All @@ -788,6 +801,7 @@ func (s *Scheduler) collectAzureForAccount(ctx context.Context, acct config.Clou
return nil, fmt.Errorf("get recommendations client: %w", err)
}
recs, err := recClient.GetAllRecommendations(ctx)
err = tolerateIncompleteSweep("azure", err)
if err != nil {
return nil, fmt.Errorf("get recommendations: %w", err)
}
Expand Down Expand Up @@ -856,6 +870,7 @@ func (s *Scheduler) collectGCPForAccount(ctx context.Context, acct config.CloudA
return nil, fmt.Errorf("get recommendations client: %w", err)
}
recs, err := recClient.GetAllRecommendations(ctx)
err = tolerateIncompleteSweep("gcp", err)
if err != nil {
return nil, fmt.Errorf("get recommendations: %w", err)
}
Expand All @@ -876,13 +891,48 @@ func (s *Scheduler) enabledAccounts(ctx context.Context, providerName string) []
return accounts
}

// tolerateIncompleteSweep converts an org-wide partial-subscription failure
// into a warning and a nil error, so the caller keeps the recommendations that
// WERE collected instead of discarding them.
//
// Every recommendation-fetch site in this file has the shape
// `recs, err := ...; if err != nil { return nil, ... }`, which throws the
// results away. That is correct for a real error, but for a partial sweep it
// would turn one flaky subscription out of fifty into a total collection
// outage -- strictly worse than the silent under-collection that
// PartialSubscriptionFailureError exists to prevent. Routing every site
// through this helper keeps the policy in one place rather than relying on
// each call site to remember it.
//
// The failed subscription IDs are named in the log so an operator can tell an
// under-collected sweep from a genuinely shrinking savings opportunity. Any
// other error is returned unchanged, preserving the existing fail-loud
// behavior.
//
// NOTE: this records the incompleteness in the log only. Surfacing it in the
// state table's last_collection_error (so the dashboard shows the sweep as
// partial) needs a partial-note threaded through collectProviderRecommendations
// and its three per-provider implementations, which is left as follow-up work.
func tolerateIncompleteSweep(providerName string, err error) error {
partial := azureprovider.AsPartialSubscriptionFailure(err)
if partial == nil {
return err
}
logging.Warnf(
"%s recommendations incomplete: %d of %d subscriptions succeeded; not queried: %s (keeping the %d that did)",
providerName, partial.Succeeded, partial.Attempted,
strings.Join(partial.FailedSubscriptionIDs(), ", "), partial.Succeeded)
return nil
}

// fetchAndConvert is a convenience for the AWS ambient path.
func (s *Scheduler) fetchAndConvert(ctx context.Context, prov provider.Provider, providerName string, accountID *string, globalCfg *config.GlobalConfig) ([]config.RecommendationRecord, error) {
recClient, err := prov.GetRecommendationsClient(ctx)
if err != nil {
return nil, fmt.Errorf("failed to get %s recommendations client: %w", providerName, err)
}
recs, err := recClient.GetAllRecommendations(ctx)
err = tolerateIncompleteSweep(providerName, err)
if err != nil {
return nil, fmt.Errorf("failed to get %s recommendations: %w", providerName, err)
}
Expand All @@ -898,6 +948,7 @@ func (s *Scheduler) fetchAndConvert(ctx context.Context, prov provider.Provider,
}
var recErr error
recs, recErr = recClient.GetRecommendations(ctx, &params)
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
Expand Down
130 changes: 130 additions & 0 deletions internal/scheduler/scheduler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,7 @@ import (
"github.com/LeanerCloud/CUDly/internal/purchase"
"github.com/LeanerCloud/CUDly/pkg/common"
"github.com/LeanerCloud/CUDly/pkg/provider"
azureprovider "github.com/LeanerCloud/CUDly/providers/azure"
"github.com/aws/aws-sdk-go-v2/aws"
"github.com/aws/aws-sdk-go-v2/service/sts"
"github.com/stretchr/testify/assert"
Expand Down Expand Up @@ -703,6 +704,105 @@ func TestSchedulerWithPurchaseManager(t *testing.T) {
assert.NotNil(t, scheduler.email)
}

// TestScheduler_FetchAndConvert_KeepsPartialSweepData is the regression guard
// for the second-order effect of returning a typed partial-failure error:
// every recommendation-fetch site in the scheduler has the shape
// `recs, err := ...; if err != nil { return nil, ... }`, which DISCARDS the
// results. Without tolerateIncompleteSweep, one flaky subscription out of
// fifty would turn a partial sweep into a total collection outage -- strictly
// worse than the silent under-collection the typed error exists to prevent.
//
// The successful subscriptions' recommendations must survive, and the call
// must not fail.
func TestScheduler_FetchAndConvert_KeepsPartialSweepData(t *testing.T) {
ctx := context.Background()

// Term/PaymentOption must be present and canonical: convertRecommendations
// deliberately drops rows with an unparseable term rather than defaulting
// one, so an under-specified fixture would pass this test for the wrong
// reason (empty in, empty out).
collected := []common.Recommendation{
{
Provider: common.ProviderAzure,
Service: common.ServiceCompute,
Account: "sub-2",
ResourceType: "Standard_D2s_v3",
Region: "westeurope",
Term: "1yr",
PaymentOption: "upfront",
Count: 3,
},
}
partial := &azureprovider.PartialSubscriptionFailureError{
Attempted: 3,
Succeeded: 2,
Failed: []azureprovider.SubscriptionFailure{
{SubscriptionID: "sub-1", Err: errors.New("throttled")},
},
}

recClient := new(MockRecommendationsClient)
recClient.On("GetAllRecommendations", mock.Anything).Return(collected, partial)
t.Cleanup(func() { recClient.AssertExpectations(t) })

prov := new(MockProvider)
prov.On("GetRecommendationsClient", mock.Anything).Return(recClient, nil)
t.Cleanup(func() { prov.AssertExpectations(t) })

s := &Scheduler{config: new(MockConfigStore)}

// globalCfg nil so the zero-results fallback branch stays out of the way;
// collected is non-empty anyway.
recs, err := s.fetchAndConvert(ctx, prov, "azure", nil, nil)

require.NoError(t, err,
"a partial multi-subscription sweep must not fail the whole collection")
require.Len(t, recs, 1,
"the subscriptions that succeeded must still be persisted, not discarded")
assert.Equal(t, 3, recs[0].Count, "the surviving subscription's recommendation must be intact")
}

// The tolerance must be narrow: any error that is NOT a partial-subscription
// failure keeps the existing fail-loud behavior.
func TestScheduler_FetchAndConvert_RealErrorStillFailsLoud(t *testing.T) {
ctx := context.Background()

recClient := new(MockRecommendationsClient)
recClient.On("GetAllRecommendations", mock.Anything).Return(nil, errors.New("credentials expired"))
t.Cleanup(func() { recClient.AssertExpectations(t) })

prov := new(MockProvider)
prov.On("GetRecommendationsClient", mock.Anything).Return(recClient, nil)
t.Cleanup(func() { prov.AssertExpectations(t) })

s := &Scheduler{config: new(MockConfigStore)}

recs, err := s.fetchAndConvert(ctx, prov, "azure", nil, nil)

require.Error(t, err, "a genuine error must still fail the collection")
assert.Contains(t, err.Error(), "credentials expired")
assert.Nil(t, recs)
}

func TestTolerateIncompleteSweep(t *testing.T) {
t.Run("partial failure is swallowed so the caller keeps its data", func(t *testing.T) {
partial := &azureprovider.PartialSubscriptionFailureError{
Attempted: 2, Succeeded: 1,
Failed: []azureprovider.SubscriptionFailure{{SubscriptionID: "sub-1", Err: errors.New("boom")}},
}
assert.NoError(t, tolerateIncompleteSweep("azure", partial))
})

t.Run("other errors pass through unchanged", func(t *testing.T) {
boom := errors.New("boom")
assert.Same(t, boom, tolerateIncompleteSweep("azure", boom))
})

t.Run("nil stays nil", func(t *testing.T) {
assert.NoError(t, tolerateIncompleteSweep("azure", nil))
})
}

// MockProvider is a mock implementation of provider.Provider.
type MockProvider struct {
mock.Mock
Expand Down Expand Up @@ -1647,6 +1747,36 @@ func TestScheduler_CollectAzureRecommendations_AllAccountsFailLoud(t *testing.T)
"the failed account must not land in SucceededAccountIDs (stale-row eviction guard)")
}

// An Azure cloud_accounts row with no azure_subscription_id must be rejected
// before a provider is built.
//
// collectAzureForAccount tags every recommendation it returns with THIS
// account's UUID, so the provider has to be pinned to this account's
// subscription. Leaving AzureSubscriptionID empty leaves the provider
// unpinned, and an unpinned provider fans out across every subscription the
// credential can see -- filing other subscriptions' recommendations under
// this account, where anyone authorized for it can read them. The row is
// misconfigured; fail loud and name the missing field rather than collecting
// data that will be attributed to the wrong account.
func TestScheduler_CollectAzureForAccount_MissingSubscriptionIDFailsLoud(t *testing.T) {
ctx := context.Background()
scheduler := &Scheduler{config: new(MockConfigStore)}

recs, err := scheduler.collectAzureForAccount(ctx, config.CloudAccount{
ID: "az-no-sub",
Provider: "azure",
AzureAuthMode: "managed_identity",
Enabled: true,
// AzureSubscriptionID deliberately empty.
})

require.Error(t, err)
assert.Nil(t, recs)
assert.Contains(t, err.Error(), "no azure_subscription_id configured",
"an unpinned Azure account must be rejected by name, not fall through to an unscoped org-wide collection")
assert.Contains(t, err.Error(), "az-no-sub", "the error must identify the misconfigured account")
}

// Test GCP recommendations with no accounts — should skip gracefully.
func TestScheduler_CollectGCPRecommendations_NoAccounts(t *testing.T) {
ctx := context.Background()
Expand Down
Loading
Loading