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
14 changes: 4 additions & 10 deletions providers/gcp/provider.go
Original file line number Diff line number Diff line change
Expand Up @@ -333,16 +333,10 @@ func (p *GCPProvider) GetAccounts(ctx context.Context) ([]common.Account, error)
}
}

// If no projects found, return at least the default project
if len(accounts) == 0 {
accounts = append(accounts, common.Account{
Provider: common.ProviderGCP,
ID: p.projectID,
Name: p.projectID,
IsDefault: true,
})
}

// Return empty if no ACTIVE projects are visible to the credentials (10-M7).
// Synthesising a fallback account from p.projectID can silently succeed
// when the account has no accessible projects, hiding auth / permission
// errors from callers. Return empty so the caller can surface the gap.
return accounts, nil
}

Expand Down
8 changes: 5 additions & 3 deletions providers/gcp/provider_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -562,6 +562,10 @@ func TestGCPProvider_GetAccounts_WithMock(t *testing.T) {
assert.Equal(t, "project-2", accounts[1].ID)
}

// TestGCPProvider_GetAccounts_Empty asserts that GetAccounts returns an empty
// slice (not a synthesised fallback project) when no ACTIVE projects are
// visible to the credentials (10-M7). A synthesised account would hide
// permission / auth errors from callers.
func TestGCPProvider_GetAccounts_Empty(t *testing.T) {
ctx := context.Background()
p := NewProviderWithProject(ctx, "default-project")
Expand All @@ -573,9 +577,7 @@ func TestGCPProvider_GetAccounts_Empty(t *testing.T) {

accounts, err := p.GetAccounts(ctx)
require.NoError(t, err)
// Should return the default project when no projects found
assert.Len(t, accounts, 1)
assert.Equal(t, "default-project", accounts[0].ID)
assert.Empty(t, accounts, "GetAccounts must return empty when no ACTIVE projects are visible (10-M7)")
}

func TestGCPProvider_GetAccounts_Error(t *testing.T) {
Expand Down
173 changes: 118 additions & 55 deletions providers/gcp/recommendations.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,9 @@ import (
"github.com/LeanerCloud/CUDly/pkg/concurrency"
"github.com/LeanerCloud/CUDly/pkg/logging"
"github.com/LeanerCloud/CUDly/providers/gcp/services/cloudsql"
"github.com/LeanerCloud/CUDly/providers/gcp/services/cloudstorage"
"github.com/LeanerCloud/CUDly/providers/gcp/services/computeengine"
"github.com/LeanerCloud/CUDly/providers/gcp/services/memorystore"
)

// defaultGCPRegionConcurrency caps the parallel per-region goroutines inside
Expand All @@ -42,13 +44,21 @@ func gcpRegionConcurrency() int {
return defaultGCPRegionConcurrency
}

// regionResult bundles the Compute Engine and Cloud SQL recommendation slices
// returned for a single GCP region. The merge in GetRecommendations walks
// regions in sorted order and appends compute then sql per region so output
// is deterministic independent of goroutine completion order.
// regionResult bundles per-service recommendation slices returned for a single
// GCP region. The merge in GetRecommendations walks regions in sorted order
// and appends compute, sql, cache, storage per region so output is
// deterministic independent of goroutine completion order.
//
// All four GCP service clients (computeengine, cloudsql, memorystore,
// cloudstorage) implement GetRecommendations and are fanned out concurrently
// when shouldIncludeService permits. Note that cache and storage purchase paths
// are advisory-only (no-op PurchaseCommitment); their recommendations are still
// surfaced so operators can see spend-optimisation signals.
type regionResult struct {
compute []common.Recommendation
sql []common.Recommendation
cache []common.Recommendation
storage []common.Recommendation
}

// RecommendationsClientAdapter aggregates GCP CUD and commitment recommendations across all services
Expand All @@ -65,9 +75,10 @@ type RecommendationsClientAdapter struct {
// - Outer: errgroup over regions, capped at gcpRegionConcurrency()
// (CUDLY_GCP_REGION_PARALLELISM, default 10) to stay polite to the
// project-scoped Recommender API quota.
// - Inner: within each region's goroutine, the (compute, cloud-sql) calls
// run as two further goroutines under a per-region sub-errgroup, so the
// per-region cost is max(compute, sql) rather than compute + sql.
// - Inner: within each region's goroutine, the four service calls
// (compute, cloud-sql, memorystore, cloudstorage) run as concurrent
// goroutines under a per-region sub-errgroup, so the per-region cost is
// max(service latencies) rather than their sum.
//
// Behaviour change vs the previous nested for-loops: per-(region, service)
// errors that were previously silently swallowed (`if err == nil { ... }`
Expand Down Expand Up @@ -132,9 +143,9 @@ func (r *RecommendationsClientAdapter) GetRecommendations(ctx context.Context, p
return nil, err
}

// Deterministic merge: walk regions in sorted order, append compute then
// sql per region. Output is stable regardless of GCP API region-list
// ordering or goroutine completion order.
// Deterministic merge: walk regions in sorted order, append compute, sql,
// cache, storage per region. Output is stable regardless of GCP API
// region-list ordering or goroutine completion order.
sortedRegions := make([]string, 0, len(results))
for region := range results {
sortedRegions = append(sortedRegions, region)
Expand All @@ -146,64 +157,109 @@ func (r *RecommendationsClientAdapter) GetRecommendations(ctx context.Context, p
res := results[region]
allRecommendations = append(allRecommendations, res.compute...)
allRecommendations = append(allRecommendations, res.sql...)
allRecommendations = append(allRecommendations, res.cache...)
allRecommendations = append(allRecommendations, res.storage...)
}
return allRecommendations, nil
}

// collectRegion fetches Compute Engine and Cloud SQL recommendations for a
// single region concurrently. Per-service errors are logged at WARN with the
// region+service tag and never propagate — the previous silent-skip-on-err
// shape is preserved (so a misconfigured project doesn't error out the whole
// recommendations refresh) but errors are now observable in logs. Extracted
// from GetRecommendations to keep that function under the gocyclo gate
// (.golangci.yml min-complexity: 15) after the post-Wait ctx.Err() block was
// added.
// collectComputeRecs fetches Compute Engine CUD recommendations for one region.
// Handles semaphore acquire/release and client construction so these branches
// are not counted toward collectRegion's cyclomatic complexity.
func (r *RecommendationsClientAdapter) collectComputeRecs(ctx context.Context, params common.RecommendationParams, region string) ([]common.Recommendation, error) {
if err := concurrency.Acquire(ctx); err != nil {
return nil, err
}
defer concurrency.Release(ctx)
client, err := computeengine.NewClient(ctx, r.projectID, region, r.clientOpts...)
if err != nil {
return nil, err
}
return client.GetRecommendations(ctx, params)
}

// collectSQLRecs fetches Cloud SQL CUD recommendations for one region.
func (r *RecommendationsClientAdapter) collectSQLRecs(ctx context.Context, params common.RecommendationParams, region string) ([]common.Recommendation, error) {
if err := concurrency.Acquire(ctx); err != nil {
return nil, err
}
defer concurrency.Release(ctx)
client, err := cloudsql.NewClient(ctx, r.projectID, region, r.clientOpts...)
if err != nil {
return nil, err
}
return client.GetRecommendations(ctx, params)
}

// collectCacheRecs fetches Memorystore recommendations for one region.
func (r *RecommendationsClientAdapter) collectCacheRecs(ctx context.Context, params common.RecommendationParams, region string) ([]common.Recommendation, error) {
if err := concurrency.Acquire(ctx); err != nil {
return nil, err
}
defer concurrency.Release(ctx)
client, err := memorystore.NewClient(ctx, r.projectID, region, r.clientOpts...)
if err != nil {
return nil, err
}
return client.GetRecommendations(ctx, params)
}

// collectStorageRecs fetches Cloud Storage recommendations for one region.
func (r *RecommendationsClientAdapter) collectStorageRecs(ctx context.Context, params common.RecommendationParams, region string) ([]common.Recommendation, error) {
if err := concurrency.Acquire(ctx); err != nil {
return nil, err
}
defer concurrency.Release(ctx)
client, err := cloudstorage.NewClient(ctx, r.projectID, region, r.clientOpts...)
if err != nil {
return nil, err
}
return client.GetRecommendations(ctx, params)
}

// collectRegion fetches recommendations for all four GCP services
// (Compute Engine, Cloud SQL, Memorystore, Cloud Storage) for a single region
// concurrently. Per-service errors are logged at WARN with the region+service
// tag and never propagate -- the previous silent-skip-on-err shape is preserved
// (so a misconfigured project doesn't error out the whole recommendations
// refresh) but errors are now observable in logs. The per-service fetch logic
// (semaphore, client construction, GetRecommendations call) is delegated to
// dedicated helpers (collectComputeRecs, collectSQLRecs, collectCacheRecs,
// collectStorageRecs) to keep this function's cyclomatic complexity under the
// gocyclo gate.
//
// Note: memorystore and cloudstorage PurchaseCommitment paths are advisory-only
// (no programmatic purchase API exists for either); their recommendations are
// surfaced so operators can see spend-optimisation signals (H-2 fix).
func (r *RecommendationsClientAdapter) collectRegion(ctx context.Context, params common.RecommendationParams, region string) regionResult {
var (
computeRecs, sqlRecs []common.Recommendation
computeErr, sqlErr error
computeRecs, sqlRecs, cacheRecs, storageRecs []common.Recommendation
computeErr, sqlErr, cacheErr, storageErr error
)

g, gctx := errgroup.WithContext(ctx)

// Per-(region, service) goroutines are leaves — they issue the actual
// Recommender API call. Acquire bounds aggregate concurrent IO across
// the whole recommendations-collection fan-out tree at the shared
// semaphore's cap (CUDLY_MAX_PARALLELISM, default 20); Release returns
// the slot. Without this bound the per-region fan-out (cap 10) ×
// per-service sub-fan-out (2) × accounts × providers can produce
// hundreds of concurrent gRPC clients that exhaust Lambda memory. If
// no semaphore is on ctx (CLI tools, unit tests), Acquire/Release are
// no-ops. See pkg/concurrency.
if shouldIncludeService(params, common.ServiceCompute) {
g.Go(func() error {
if err := concurrency.Acquire(gctx); err != nil {
computeErr = err
return nil
}
defer concurrency.Release(gctx)
client, err := computeengine.NewClient(gctx, r.projectID, region, r.clientOpts...)
if err != nil {
computeErr = err
return nil
}
computeRecs, computeErr = client.GetRecommendations(gctx, params)
computeRecs, computeErr = r.collectComputeRecs(gctx, params, region)
return nil
})
}
if shouldIncludeService(params, common.ServiceRelationalDB) {
g.Go(func() error {
if err := concurrency.Acquire(gctx); err != nil {
sqlErr = err
return nil
}
defer concurrency.Release(gctx)
client, err := cloudsql.NewClient(gctx, r.projectID, region, r.clientOpts...)
if err != nil {
sqlErr = err
return nil
}
sqlRecs, sqlErr = client.GetRecommendations(gctx, params)
sqlRecs, sqlErr = r.collectSQLRecs(gctx, params, region)
return nil
})
}
if shouldIncludeService(params, common.ServiceCache) {
g.Go(func() error {
cacheRecs, cacheErr = r.collectCacheRecs(gctx, params, region)
return nil
})
}
if shouldIncludeService(params, common.ServiceStorage) {
g.Go(func() error {
storageRecs, storageErr = r.collectStorageRecs(gctx, params, region)
return nil
})
}
Expand All @@ -215,8 +271,14 @@ func (r *RecommendationsClientAdapter) collectRegion(ctx context.Context, params
if sqlErr != nil {
logging.Warnf("GCP %s cloudsql recommendations: %v", region, sqlErr)
}
if cacheErr != nil {
logging.Warnf("GCP %s memorystore recommendations: %v", region, cacheErr)
}
if storageErr != nil {
logging.Warnf("GCP %s cloudstorage recommendations: %v", region, storageErr)
}

return regionResult{compute: computeRecs, sql: sqlRecs}
return regionResult{compute: computeRecs, sql: sqlRecs, cache: cacheRecs, storage: storageRecs}
}

// GetRecommendationsForService retrieves GCP commitment recommendations for a specific service
Expand All @@ -235,10 +297,11 @@ func (r *RecommendationsClientAdapter) GetAllRecommendations(ctx context.Context

// getRegions retrieves available GCP regions for the project
func (r *RecommendationsClientAdapter) getRegions(ctx context.Context) ([]string, error) {
// Create a temporary provider to get regions
provider := NewProviderWithProject(ctx, r.projectID, r.clientOpts...)
// Create a temporary provider to get regions. The local variable is named
// p (not provider) to avoid shadowing the imported provider package (10-N3).
p := NewProviderWithProject(ctx, r.projectID, r.clientOpts...)

regions, err := provider.GetRegions(ctx)
regions, err := p.GetRegions(ctx)
if err != nil {
return nil, err
}
Expand Down
75 changes: 75 additions & 0 deletions providers/gcp/recommendations_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -160,6 +160,81 @@ func TestRecommendationsClientAdapter_GetRecommendations_PropagatesContextCancel
"GetRecommendations must propagate the parent ctx error")
}

// TestRegionResult_HasCacheAndStorageFields is a compile-time regression test
// for H-2 (GCP broad audit): regionResult must carry cache and storage slices
// so collectRegion can fan out to memorystore and cloudstorage in addition to
// compute and cloudsql. If this test stops compiling, the wiring was reverted.
func TestRegionResult_HasCacheAndStorageFields(t *testing.T) {
recs := []common.Recommendation{{Provider: common.ProviderGCP}}

// Field access verifies that regionResult has the full four-service shape.
result := regionResult{
compute: recs,
sql: recs,
cache: recs,
storage: recs,
}

assert.Len(t, result.compute, 1, "compute field must be present on regionResult")
assert.Len(t, result.sql, 1, "sql field must be present on regionResult")
assert.Len(t, result.cache, 1, "cache field must be present on regionResult (H-2)")
assert.Len(t, result.storage, 1, "storage field must be present on regionResult (H-2)")
}

// TestShouldIncludeService_Cache_Storage verifies that shouldIncludeService
// correctly routes ServiceCache and ServiceStorage requests, which is a
// prerequisite for H-2 (wiring them into collectRegion).
func TestShouldIncludeService_Cache_Storage(t *testing.T) {
tests := []struct {
name string
params common.RecommendationParams
service common.ServiceType
expected bool
}{
{
name: "all services includes Cache",
params: common.RecommendationParams{},
service: common.ServiceCache,
expected: true,
},
{
name: "all services includes Storage",
params: common.RecommendationParams{},
service: common.ServiceStorage,
expected: true,
},
{
name: "Cache-scoped request includes Cache",
params: common.RecommendationParams{Service: common.ServiceCache},
service: common.ServiceCache,
expected: true,
},
{
name: "Cache-scoped request excludes Storage",
params: common.RecommendationParams{Service: common.ServiceCache},
service: common.ServiceStorage,
expected: false,
},
{
name: "Storage-scoped request includes Storage",
params: common.RecommendationParams{Service: common.ServiceStorage},
service: common.ServiceStorage,
expected: true,
},
{
name: "Storage-scoped request excludes Compute",
params: common.RecommendationParams{Service: common.ServiceStorage},
service: common.ServiceCompute,
expected: false,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
assert.Equal(t, tt.expected, shouldIncludeService(tt.params, tt.service))
})
}
}

// TestGCPRegionConcurrency pins the env-knob parsing for
// CUDLY_GCP_REGION_PARALLELISM.
func TestGCPRegionConcurrency(t *testing.T) {
Expand Down
Loading
Loading