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
10 changes: 7 additions & 3 deletions internal/config/store_postgres.go
Original file line number Diff line number Diff line change
Expand Up @@ -80,7 +80,8 @@ func getGlobalConfigFrom(ctx context.Context, q globalConfigExecutor) (*GlobalCo
grace_period_days,
recommendations_cache_stale_hours, recommendations_lookback_days,
COALESCE(purchase_delay_hours, 0),
COALESCE(laddering_enabled, false)
COALESCE(laddering_enabled, false),
COALESCE(ladder_execution_enabled, false)
FROM global_config
WHERE id = 1
`
Expand Down Expand Up @@ -111,6 +112,7 @@ func getGlobalConfigFrom(ctx context.Context, q globalConfigExecutor) (*GlobalCo
&config.RecommendationsLookbackDays,
&config.PurchaseDelayHours,
&config.LadderingEnabled,
&config.LadderExecutionEnabled,
)

if err != nil {
Expand Down Expand Up @@ -209,8 +211,8 @@ func saveGlobalConfigWith(ctx context.Context, q globalConfigExecutor, config *G
auto_collect, collection_schedule, notification_days_before,
grace_period_days,
recommendations_cache_stale_hours, recommendations_lookback_days,
purchase_delay_hours, laddering_enabled
) VALUES (1, $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21)
purchase_delay_hours, laddering_enabled, ladder_execution_enabled
) VALUES (1, $1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22)
ON CONFLICT (id) DO UPDATE SET
enabled_providers = $1,
notification_email = $2,
Expand All @@ -233,6 +235,7 @@ func saveGlobalConfigWith(ctx context.Context, q globalConfigExecutor, config *G
recommendations_lookback_days = $19,
purchase_delay_hours = $20,
laddering_enabled = $21,
ladder_execution_enabled = $22,
updated_at = NOW()
`

Expand Down Expand Up @@ -290,6 +293,7 @@ func saveGlobalConfigWith(ctx context.Context, q globalConfigExecutor, config *G
recommendationsLookbackDays,
config.PurchaseDelayHours,
config.LadderingEnabled,
config.LadderExecutionEnabled,
)

if err != nil {
Expand Down
9 changes: 8 additions & 1 deletion internal/config/store_postgres_pgxmock_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,7 @@ func TestPGXMock_GetGlobalConfig_Success(t *testing.T) {
"recommendations_cache_stale_hours", "recommendations_lookback_days",
"purchase_delay_hours",
"laddering_enabled",
"ladder_execution_enabled",
}
rows := pgxmock.NewRows(cols).AddRow(
[]string{"aws"}, strPtr("ops@example.com"), true,
Expand All @@ -71,6 +72,7 @@ func TestPGXMock_GetGlobalConfig_Success(t *testing.T) {
24, 7,
0,
false,
false,
)
mock.ExpectQuery("SELECT").WillReturnRows(rows)

Expand Down Expand Up @@ -112,6 +114,7 @@ func TestPGXMock_GetGlobalConfig_GracePeriodDays(t *testing.T) {
"recommendations_cache_stale_hours", "recommendations_lookback_days",
"purchase_delay_hours",
"laddering_enabled",
"ladder_execution_enabled",
}
baseRow := func(graceJSON string) []any {
return []any{
Expand All @@ -124,6 +127,7 @@ func TestPGXMock_GetGlobalConfig_GracePeriodDays(t *testing.T) {
24, 7,
0,
false,
false,
}
}

Expand Down Expand Up @@ -181,6 +185,7 @@ var globalConfigCols = []string{
"recommendations_cache_stale_hours", "recommendations_lookback_days",
"purchase_delay_hours",
"laddering_enabled",
"ladder_execution_enabled",
}

// TestPGXMock_UpdateGlobalConfigAtomic_LockedReadModifyWrite proves the F2
Expand Down Expand Up @@ -209,6 +214,7 @@ func TestPGXMock_UpdateGlobalConfigAtomic_LockedReadModifyWrite(t *testing.T) {
24, 7,
48,
false, // laddering_enabled = false
false, // ladder_execution_enabled = false
)

// Strict order: the SELECT and the UPSERT must sit between the same
Expand All @@ -217,7 +223,7 @@ func TestPGXMock_UpdateGlobalConfigAtomic_LockedReadModifyWrite(t *testing.T) {
mock.ExpectExec("pg_advisory_xact_lock").WithArgs(pgxmock.AnyArg()).
WillReturnResult(pgxmock.NewResult("SELECT", 1))
mock.ExpectQuery("FROM global_config").WillReturnRows(seeded)
mock.ExpectExec("INSERT INTO global_config").WithArgs(anyArgsCfg(21)...).
mock.ExpectExec("INSERT INTO global_config").WithArgs(anyArgsCfg(22)...).
WillReturnResult(pgxmock.NewResult("INSERT", 1))
mock.ExpectCommit()

Expand Down Expand Up @@ -263,6 +269,7 @@ func TestPGXMock_UpdateGlobalConfigAtomic_ApplyErrorRollsBack(t *testing.T) {
24, 7,
0,
false,
false,
)

mock.ExpectBegin()
Expand Down
9 changes: 9 additions & 0 deletions internal/config/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -69,6 +69,15 @@ type GlobalConfig struct {
// engine runs fire regardless of per-account LadderConfig.Enabled settings.
// Set to true to allow per-account configs to activate individually.
LadderingEnabled bool `json:"laddering_enabled" db:"laddering_enabled"`

// LadderExecutionEnabled gates the write side of the ladder capability
// (migration 000083). BOTH LadderingEnabled AND LadderExecutionEnabled
// must be true for PurchaseLayer / ReshapeBuffer to be wired with real
// AWS SDK clients. Default false: existing deployments that enable
// laddering produce plans but never call AWS purchase APIs until an
// operator explicitly opts in. Fail-loud: wireLadderWriteSide returns
// a typed ErrLadderExecutionDisabled when this is false.
LadderExecutionEnabled bool `json:"ladder_execution_enabled" db:"ladder_execution_enabled"`
}

// DefaultGracePeriodDays is the fallback window used when a provider
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
ALTER TABLE global_config
DROP COLUMN IF EXISTS ladder_execution_enabled;
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
-- Migration 000083: ladder_execution_enabled global kill-switch.
--
-- Stacks on laddering_enabled (migration 000079): BOTH must be true for
-- the ladder capability write side (PurchaseLayer / ReshapeBuffer) to be
-- wired with real AWS clients. Default FALSE means existing deployments
-- that enable laddering produce plans but never call AWS purchase APIs until
-- an operator explicitly opts in (fail-loud, no silent fallback).
--
-- Idempotent: ADD COLUMN IF NOT EXISTS.

ALTER TABLE global_config
ADD COLUMN IF NOT EXISTS ladder_execution_enabled BOOLEAN NOT NULL DEFAULT FALSE;
16 changes: 7 additions & 9 deletions internal/server/handler_ladder.go
Original file line number Diff line number Diff line change
Expand Up @@ -123,7 +123,7 @@ func (app *Application) handleLadderRun(ctx context.Context) (*LadderRunResult,
}

now := time.Now().UTC()
result := app.runLadderConfigs(ctx, allConfigs, ownAccountID, region, term, paymentOpt, now)
result := app.runLadderConfigs(ctx, allConfigs, ownAccountID, region, term, paymentOpt, now, globalCfg.LadderExecutionEnabled)

log.Printf("ladder_run done: planned=%d skipped_cadence=%d skipped_disabled=%d skipped_multi_account=%d errored=%d",
result.Planned, result.SkippedCadence, result.SkippedDisabled, result.SkippedMultiAccount, result.Errored)
Expand Down Expand Up @@ -192,10 +192,11 @@ func (app *Application) runLadderConfigs(
term pkgladder.Term,
paymentOpt pkgladder.PaymentOption,
now time.Time,
executionEnabled bool,
) *LadderRunResult {
result := &LadderRunResult{}
for i := range configs {
result.record(app.processOneLadderConfig(ctx, &configs[i], ownAccountID, region, term, paymentOpt, now))
result.record(app.processOneLadderConfig(ctx, &configs[i], ownAccountID, region, term, paymentOpt, now, executionEnabled))
}
return result
}
Expand All @@ -211,6 +212,7 @@ func (app *Application) processOneLadderConfig(
term pkgladder.Term,
paymentOpt pkgladder.PaymentOption,
now time.Time,
executionEnabled bool,
) ladderConfigOutcome {
if !dbCfg.Enabled {
log.Printf("ladder_run: config %s: enabled=false, skipping", dbCfg.ID)
Expand Down Expand Up @@ -248,14 +250,10 @@ func (app *Application) processOneLadderConfig(
return outcomeSkippedCadence
}

// Build the LadderCapability for this account.
if app.LadderCapabilityFactory == nil {
log.Printf("ladder_run: config %s: LadderCapabilityFactory is nil (not wired), erroring", dbCfg.ID)
return outcomeErrored
}
capability, err := app.LadderCapabilityFactory(ctx, region, cloudAcct.ExternalID)
// Build and wire the LadderCapability for this account.
capability, err := app.buildAndWireCapability(ctx, region, cloudAcct.ExternalID, executionEnabled)
if err != nil {
log.Printf("ladder_run: config %s: failed to build ladder capability: %v", dbCfg.ID, err)
log.Printf("ladder_run: config %s: %v", dbCfg.ID, err)
return outcomeErrored
}

Expand Down
6 changes: 3 additions & 3 deletions internal/server/handler_ladder_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -467,7 +467,7 @@ func TestHandleLadderRun_MultiAccountSkip_CountedAndIsolated(t *testing.T) {

// Put the foreign config first to prove isolation is order-independent.
configs := []config.LadderConfigDB{cfgForeign, cfgHealthy}
result := app.runLadderConfigs(ctx, configs, ownAccount, "us-east-1", pkgladder.Term1Year, pkgladder.PaymentNoUpfront, now)
result := app.runLadderConfigs(ctx, configs, ownAccount, "us-east-1", pkgladder.Term1Year, pkgladder.PaymentNoUpfront, now, false)

require.NotNil(t, result)
// (a) The foreign config must be counted as SkippedMultiAccount.
Expand Down Expand Up @@ -869,7 +869,7 @@ func TestProcessOneLadderConfig_CadenceDBError_Errored(t *testing.T) {
},
}

result := app.runLadderConfigs(ctx, []config.LadderConfigDB{dbCfg}, ownAccount, "us-east-1", pkgladder.Term1Year, pkgladder.PaymentNoUpfront, now)
result := app.runLadderConfigs(ctx, []config.LadderConfigDB{dbCfg}, ownAccount, "us-east-1", pkgladder.Term1Year, pkgladder.PaymentNoUpfront, now, false)

assert.Equal(t, 1, result.Errored, "a cadence lookup error must count the config Errored")
assert.Equal(t, 0, result.Planned)
Expand Down Expand Up @@ -993,7 +993,7 @@ func TestHandleLadderRun_MultiConfigIsolation(t *testing.T) {
// Order the broken config first to prove a leading failure does not abort
// the healthy config that follows.
configs := []config.LadderConfigDB{cfgBroken, cfgHealthy}
result := app.runLadderConfigs(ctx, configs, ownAccount, "us-east-1", pkgladder.Term1Year, pkgladder.PaymentNoUpfront, now)
result := app.runLadderConfigs(ctx, configs, ownAccount, "us-east-1", pkgladder.Term1Year, pkgladder.PaymentNoUpfront, now, false)

require.NotNil(t, result)
assert.Equal(t, 1, result.Planned, "the healthy config must still be planned despite the broken one")
Expand Down
143 changes: 143 additions & 0 deletions internal/server/ladder_write.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,143 @@
package server

import (
"context"
"errors"
"fmt"

awsconfig "github.com/aws/aws-sdk-go-v2/config"

"github.com/LeanerCloud/CUDly/pkg/exchange"
pkgladder "github.com/LeanerCloud/CUDly/pkg/ladder"
awsprovider "github.com/LeanerCloud/CUDly/providers/aws"
awsladder "github.com/LeanerCloud/CUDly/providers/aws/ladder"
ec2svc "github.com/LeanerCloud/CUDly/providers/aws/services/ec2"
)

// exchangeRunnerAdapter bridges internal/server wiring (exchange store, EC2 exchange
// client, offering lookup) with the exchangeRunner seam expected by AWSLadder. It
// satisfies the unexported providers/aws/ladder.exchangeRunner interface via Go
// structural typing: the concrete RunAutoExchange method signature matches the
// interface definition, so the compiler accepts this type wherever exchangeRunner
// is expected without the caller naming the interface.
//
// The adapter owns the full client construction and conversion that
// executeRIExchangeReshape performs for the standalone RI-exchange path, adapted
// for ladder runs: LadderRunID and DryRun are forwarded from the seam arguments
// directly into RunAutoExchangeParams.
type exchangeRunnerAdapter struct {
app *Application
region string
accountID string
}

// RunAutoExchange implements providers/aws/ladder.exchangeRunner. It constructs
// fresh AWS clients, lists convertible RIs and utilization, converts them for the
// exchange package, then delegates to exchange.RunAutoExchange.
//
// ladderRunID is forwarded to RunAutoExchangeParams.LadderRunID so exchange scopes
// its pending-cancellation to the ladder origin (issue #1348 / gap G10). DryRun
// is forwarded so the exchange engine skips all mutations and returns Simulated
// outcomes when the ladder run was started in dry-run mode.
func (a *exchangeRunnerAdapter) RunAutoExchange(ctx context.Context, cfg exchange.RIExchangeConfig, ladderRunID *string, dryRun bool) (*exchange.AutoExchangeResult, error) {
awsCfg, err := awsconfig.LoadDefaultConfig(ctx, awsconfig.WithRegion(a.region))
if err != nil {
return nil, fmt.Errorf("exchangeRunnerAdapter: load AWS config: %w", err)
}

ec2Client := awsprovider.NewEC2ClientDirect(awsCfg)
recsClient := awsprovider.NewRecommendationsClientDirect(awsCfg)

instances, err := ec2Client.ListConvertibleReservedInstances(ctx)
if err != nil {
return nil, fmt.Errorf("exchangeRunnerAdapter: list convertible RIs: %w", err)
}
utilData, err := recsClient.GetRIUtilization(ctx, cfg.LookbackDays)
if err != nil {
return nil, fmt.Errorf("exchangeRunnerAdapter: get RI utilization: %w", err)
}

riInfos, utilInfos, riMetadata := convertForAutoExchange(instances, utilData)
store := newConfigExchangeStoreAdapter(a.app.Config)

lookupOffering := func(ctx context.Context, instanceType, productDesc, tenancy, scope string, duration int64) (string, error) {
return ec2Client.FindConvertibleOffering(ctx, ec2svc.FindConvertibleOfferingParams{
InstanceType: instanceType,
ProductDescription: productDesc,
Tenancy: tenancy,
Scope: scope,
Duration: duration,
})
}

return exchange.RunAutoExchange(ctx, exchange.RunAutoExchangeParams{
Store: store,
ExchangeClient: exchange.NewExchangeClient(awsCfg),
LookupOffering: lookupOffering,
RIs: riInfos,
Utilization: utilInfos,
Config: cfg,
AccountID: a.accountID,
Region: a.region,
DashboardURL: a.app.appConfig.DashboardURL,
RIMetadata: riMetadata,
LadderRunID: ladderRunID,
DryRun: dryRun,
})
}

// buildAndWireCapability constructs a LadderCapability via the factory and wires its
// write side. Extracted from processOneLadderConfig to keep that function's cyclomatic
// complexity below the project threshold (10). The returned error carries no config
// ID: the sole caller already prefixes its log line with the config ID, so repeating
// it here would duplicate it in the output.
func (app *Application) buildAndWireCapability(ctx context.Context, region, accountID string, executionEnabled bool) (pkgladder.LadderCapability, error) {
if app.LadderCapabilityFactory == nil {
return nil, errors.New("LadderCapabilityFactory is nil (not wired)")
}
capability, err := app.LadderCapabilityFactory(ctx, region, accountID)
if err != nil {
return nil, fmt.Errorf("failed to build ladder capability: %w", err)
}
return app.wireLadderWriteSide(ctx, executionEnabled, region, accountID, capability)
}

// wireLadderWriteSide wires the write side of a LadderCapability if and only if
// cap is a *awsladder.AWSLadder. For test fakes (non-AWSLadder implementations)
// it returns cap unchanged so existing handler tests continue to work without
// modifying their fake capabilities.
//
// When executionEnabled is false, WireWriteSideDisabled is called: the ladder
// accepts PurchaseLayer / ReshapeBuffer calls but immediately returns
// ErrLadderExecutionDisabled without touching any AWS API. When executionEnabled
// is true, WireWriteSide wires real EC2 and Savings Plans clients plus the
// exchangeRunnerAdapter that forwards ladderRunID and dryRun to the exchange
// package's seam.
func (app *Application) wireLadderWriteSide(
ctx context.Context,
executionEnabled bool,
region, accountID string,
capability pkgladder.LadderCapability,
) (pkgladder.LadderCapability, error) {
l, ok := capability.(*awsladder.AWSLadder)
if !ok {
// Test fake or non-AWS capability: plan-only invariant holds without wiring.
return capability, nil
}
if !executionEnabled {
wired, err := awsladder.WireWriteSideDisabled(l)
if err != nil {
return nil, fmt.Errorf("wireLadderWriteSide: disabled: %w", err)
}
return wired, nil
}
awsCfg, err := awsconfig.LoadDefaultConfig(ctx, awsconfig.WithRegion(region))
if err != nil {
return nil, fmt.Errorf("wireLadderWriteSide: load AWS config: %w", err)
}
wired, err := awsladder.WireWriteSide(l, awsCfg, &exchangeRunnerAdapter{app: app, region: region, accountID: accountID})
if err != nil {
return nil, fmt.Errorf("wireLadderWriteSide: %w", err)
}
return wired, nil
}
Loading
Loading