From a2e2465fcb22b893b2a08e8155d7fa6697512d4e Mon Sep 17 00:00:00 2001 From: Cristian Magherusan-Stanciu Date: Wed, 19 Aug 2026 12:22:16 +0200 Subject: [PATCH 1/2] fix(plans): make the ramp advance idempotent per step (#1669) IncrementPlanCurrentStep advanced the ramp with a blind CurrentStep++ that carried no notion of WHICH step was completing, and its caller could not supply one either. A multi-account plan produces one execution per account per ramp step, so retrying two separately-failed accounts of a single step advanced the ramp twice: the plan reported itself a step further along than the commitment it had actually bought, and a later step's purchase silently dropped out of its accounting. Measured against a real Postgres before the fix, for a 4-step plan on step 2 whose step-3 fan-out committed 1 of 3 accounts: the partial root correctly did not advance (step 2), the first per-account retry advanced to 3, and the second advanced to 4. At that point IsComplete() reported the ramp finished and next_execution_date was cleared, after only 3 of 4 steps had been bought. Replace it with CompletePlanStep(ctx, planID, stepNumber), which advances only when the completing step is exactly CurrentStep+1, inside the existing FOR UPDATE transaction. Completing a step at or below CurrentStep is a no-op, so a second per-account retry cannot advance the ramp again. Completing a step further ahead is refused rather than jumping, because advancing over steps that never completed would overstate what the plan has bought. The step comes from the execution's own step_number column (NOT NULL DEFAULT 1 since the initial schema), never from the plan state that is itself in doubt; a plan-attributed row with a non-positive step is reported instead of guessed. purchase.getOrCreateExecution stamped step_number with the COUNT of completed steps rather than the 1-based step being executed, disagreeing with api.createPurchaseExecutionsTx (CurrentStep + i + 1). That off-by-one was invisible under a blind increment but would have frozen the ramp once the advance keyed on the value, so it is corrected here, and migration 000098 retargets the executable rows already carrying the old convention. Without that backfill every in-flight plan execution would silently stop advancing its ramp on deploy, which is the same under-buy this change exists to fix. Terminal rows are left alone: their step_number is audit trail and no discriminator separates an old-convention one from a correctly stamped one whose ramp has since moved past it. The backfill also skips any row whose plan and step already has a successful sibling: that shape is a per-account retry successor for a step another account's retry already counted, which is correctly stamped, and retargeting it would make it buy its own step's tranche while advancing the ramp past the next one. Two limits stay open and are unchanged by this commit. A step still counts as completed when one execution for it runs clean, not when every account of a multi-account plan has bought. And a step that never completes now freezes the ramp rather than letting later steps advance over it, which is the safer of the two wrong answers but leaves no durable evidence beyond a log line. Regression coverage runs the real Manager against a real migrated Postgres (internal/purchase/ramp_step_progress_integration_test.go): the three-account fan-out with two failures and two retries, which fails pre-fix by assertion with CurrentStep=4; an 8-way concurrent completion of one step; and a forced interleaving where an outside transaction holds the plan row lock and commits step 3 while a concurrent completion of the same step is blocked on it, proving the blocked call observes the committed step and no-ops. Closes #1669 --- internal/analytics/collector_test.go | 2 +- .../handler_purchases_retry_fanout_test.go | 2 +- internal/api/handler_purchases_test.go | 4 +- internal/config/interfaces.go | 14 +- internal/config/store_postgres.go | 54 ++- .../store_postgres_complete_step_test.go | 278 +++++++++++ .../store_postgres_increment_step_test.go | 211 --------- ...98_backfill_execution_step_number.down.sql | 13 + ...0098_backfill_execution_step_number.up.sql | 65 +++ ...098_backfill_execution_step_number_test.go | 170 +++++++ internal/mocks/stores.go | 8 +- internal/purchase/approvals_test.go | 17 +- internal/purchase/armed_redrive_test.go | 2 +- internal/purchase/coverage_extra_test.go | 9 +- internal/purchase/execution.go | 26 +- internal/purchase/execution_test.go | 68 ++- internal/purchase/manager.go | 2 +- internal/purchase/manager_test.go | 12 +- .../purchase/money_path_regression_test.go | 4 +- internal/purchase/notifications.go | 13 +- internal/purchase/notifications_test.go | 7 +- .../ramp_step_progress_integration_test.go | 443 ++++++++++++++++++ internal/server/test_helpers_test.go | 2 +- 23 files changed, 1157 insertions(+), 269 deletions(-) create mode 100644 internal/config/store_postgres_complete_step_test.go delete mode 100644 internal/config/store_postgres_increment_step_test.go create mode 100644 internal/database/postgres/migrations/000098_backfill_execution_step_number.down.sql create mode 100644 internal/database/postgres/migrations/000098_backfill_execution_step_number.up.sql create mode 100644 internal/database/postgres/migrations/000098_backfill_execution_step_number_test.go create mode 100644 internal/purchase/ramp_step_progress_integration_test.go diff --git a/internal/analytics/collector_test.go b/internal/analytics/collector_test.go index 4f19e8583..30baf7544 100644 --- a/internal/analytics/collector_test.go +++ b/internal/analytics/collector_test.go @@ -160,7 +160,7 @@ func (m *mockConfigStore) UpdatePurchasePlan(ctx context.Context, plan *config.P return nil } -func (m *mockConfigStore) IncrementPlanCurrentStep(ctx context.Context, planID string) error { +func (m *mockConfigStore) CompletePlanStep(ctx context.Context, planID string, stepNumber int) error { return nil } diff --git a/internal/api/handler_purchases_retry_fanout_test.go b/internal/api/handler_purchases_retry_fanout_test.go index 480776be4..a6f29ec86 100644 --- a/internal/api/handler_purchases_retry_fanout_test.go +++ b/internal/api/handler_purchases_retry_fanout_test.go @@ -171,7 +171,7 @@ func executeSuccessorForReal(t *testing.T, successor *config.PurchaseExecution) // Plan progress advances only on a fully clean run; .Maybe() so a run that // ends partial/failed does not turn into a mock-expectation failure that // would mask the purchase-count assertion the caller actually makes. - store.On("IncrementPlanCurrentStep", mock.Anything, fanoutPlanID).Return(nil).Maybe() + store.On("CompletePlanStep", mock.Anything, fanoutPlanID, mock.Anything).Return(nil).Maybe() store.GetPurchasePlanFn = func(_ context.Context, planID string) (*config.PurchasePlan, error) { return &config.PurchasePlan{ID: planID, Name: "Plan 1537"}, nil diff --git a/internal/api/handler_purchases_test.go b/internal/api/handler_purchases_test.go index d54e72f68..bc3075325 100644 --- a/internal/api/handler_purchases_test.go +++ b/internal/api/handler_purchases_test.go @@ -5653,14 +5653,14 @@ func TestHandler_approvePurchaseViaSession_FourEyesOn_DifferentApproverSucceeds( Status: "pending", CreatedByUserID: &creatorID, } - approved := &config.PurchaseExecution{ExecutionID: execID, PlanID: planID, Status: "approved"} + approved := &config.PurchaseExecution{ExecutionID: execID, PlanID: planID, Status: "approved", StepNumber: 1} mockConfig.On("GetExecutionByID", ctx, execID).Return(exec, nil) mockConfig.On("GetGlobalConfig", ctx).Return(fourEyesCfgOn(), nil) mockConfig.On("TransitionExecutionStatus", ctx, execID, []string{"pending", "notified"}, "approved", &approverID).Return(approved, nil) plan := &config.PurchasePlan{ID: planID, Name: "test-plan"} mockConfig.On("GetPurchasePlan", ctx, planID).Return(plan, nil) mockConfig.On("SavePurchaseExecution", ctx, mock.AnythingOfType("*config.PurchaseExecution")).Return(nil) - mockConfig.On("IncrementPlanCurrentStep", ctx, planID).Return(nil) + mockConfig.On("CompletePlanStep", ctx, planID, 1).Return(nil) mockAuth := new(MockAuthService) mockAuth.On("ValidateSession", ctx, "sess-tok").Return(&Session{UserID: approverID, Email: approverEmail}, nil) diff --git a/internal/config/interfaces.go b/internal/config/interfaces.go index b20d06c25..8f5adea63 100644 --- a/internal/config/interfaces.go +++ b/internal/config/interfaces.go @@ -30,10 +30,16 @@ type StoreInterface interface { CreatePurchasePlan(ctx context.Context, plan *PurchasePlan) error GetPurchasePlan(ctx context.Context, planID string) (*PurchasePlan, error) UpdatePurchasePlan(ctx context.Context, plan *PurchasePlan) error - // IncrementPlanCurrentStep atomically advances the ramp schedule for planID - // inside a SELECT FOR UPDATE transaction, preventing the concurrent-write - // lost-update race described in issue #1071. - IncrementPlanCurrentStep(ctx context.Context, planID string) error + // CompletePlanStep records that ramp step stepNumber finished and advances + // the plan's schedule to it, inside a SELECT FOR UPDATE transaction that + // prevents the concurrent-write lost-update race of issue #1071. + // + // It is idempotent in stepNumber: completing a step the plan has already + // counted is a no-op. That is what keeps a multi-account ramp step from + // advancing twice when an operator retries two separately-failed accounts + // of the same step (issue #1669). A step still counts as completed when one + // execution for it ran clean, not when every account has bought. + CompletePlanStep(ctx context.Context, planID string, stepNumber int) error // UpdatePurchasePlanTx is the tx-accepting variant of UpdatePurchasePlan. // Used from createPlannedPurchases' WithTx block so the per-row // SavePurchaseExecutionTx writes and the plan's next_execution_date diff --git a/internal/config/store_postgres.go b/internal/config/store_postgres.go index a2849b4e6..3fce71f18 100644 --- a/internal/config/store_postgres.go +++ b/internal/config/store_postgres.go @@ -548,7 +548,7 @@ const purchasePlanSelectCols = ` // scanPurchasePlanRow deserialises one purchase_plans row returned by QueryRow // or the first row of a query with FOR UPDATE. Extracted so GetPurchasePlan -// and IncrementPlanCurrentStep share the same scan logic. +// and CompletePlanStep share the same scan logic. func scanPurchasePlanRow(row pgx.Row) (*PurchasePlan, error) { var plan PurchasePlan var servicesJSON, rampScheduleJSON []byte @@ -604,13 +604,34 @@ func (s *PostgresStore) GetPurchasePlan(ctx context.Context, planID string) (*Pu return plan, nil } -// IncrementPlanCurrentStep atomically advances the ramp schedule for planID -// inside a transaction. The row is locked with SELECT FOR UPDATE so concurrent -// callers (overlapping Lambda invocations, multi-tick cron) cannot both read -// the same CurrentStep value and both write CurrentStep+1, skipping a step. +// CompletePlanStep records that ramp step stepNumber of planID completed and +// advances the schedule to it, inside a transaction. The row is locked with +// SELECT FOR UPDATE so concurrent callers (overlapping Lambda invocations, +// multi-tick cron) cannot both read the same CurrentStep value and both write +// CurrentStep+1, skipping a step (issue #1071). +// +// The operation is idempotent in stepNumber rather than a blind increment: a +// multi-account plan produces one execution per account per ramp step, and +// retrying two separately-failed accounts of the same step used to advance the +// ramp twice, so the plan reported itself a step further along than the +// commitment it had actually bought (issue #1669). Completing a step at or +// below CurrentStep is therefore a no-op. +// +// Completing a step more than one beyond CurrentStep is refused with an error +// instead of jumping: the steps in between never completed, and advancing over +// them would silently overstate how much the plan has bought. +// +// Granularity: "the step completed" still means at least one execution for that +// step ran clean, not that every account of a multi-account plan bought. That +// was true before #1669 too and is unchanged here; #1669 only stops one step +// being counted more than once. +// // Returns nil when the plan no longer exists (deleted between execution and // progress update) so the caller is not penalized for a race it cannot control. -func (s *PostgresStore) IncrementPlanCurrentStep(ctx context.Context, planID string) error { +func (s *PostgresStore) CompletePlanStep(ctx context.Context, planID string, stepNumber int) error { + if stepNumber <= 0 { + return fmt.Errorf("cannot complete ramp step %d of plan %s: step numbers are 1-based", stepNumber, planID) + } return s.WithTx(ctx, func(tx pgx.Tx) error { row := tx.QueryRow(ctx, purchasePlanSelectCols+` WHERE id = $1 FOR UPDATE`, planID) plan, err := scanPurchasePlanRow(row) @@ -621,8 +642,23 @@ func (s *PostgresStore) IncrementPlanCurrentStep(ctx context.Context, planID str return fmt.Errorf("failed to lock purchase plan %s: %w", planID, err) } - if !plan.RampSchedule.IsComplete() { - plan.RampSchedule.CurrentStep++ + switch { + case plan.RampSchedule.IsComplete(): + // The ramp finished; a trailing execution cannot extend it. Fall + // through to the write anyway so a completed ramp always ends up + // with next_execution_date cleared (the invariant issue #1071 + // established). + case plan.RampSchedule.CurrentStep >= stepNumber: + // Already counted: a second successful execution of a step the ramp + // has passed, which is what a per-account retry of a + // partially-failed fan-out produces. Nothing to advance. + return nil + case plan.RampSchedule.CurrentStep != stepNumber-1: + return fmt.Errorf("refusing to advance plan %s from ramp step %d to %d: step(s) %d-%d never completed", + planID, plan.RampSchedule.CurrentStep, stepNumber, + plan.RampSchedule.CurrentStep+1, stepNumber-1) + default: + plan.RampSchedule.CurrentStep = stepNumber } if !plan.RampSchedule.IsComplete() { @@ -634,7 +670,7 @@ func (s *PostgresStore) IncrementPlanCurrentStep(ctx context.Context, planID str now := time.Now() plan.LastExecutionDate = &now - // Refresh updated_at on every increment. The plan was read from the DB + // Refresh updated_at on every advance. The plan was read from the DB // with its previous UpdatedAt, so UpdatePurchasePlanTx's zero-value // guard would otherwise persist a stale updated_at timestamp. plan.UpdatedAt = now diff --git a/internal/config/store_postgres_complete_step_test.go b/internal/config/store_postgres_complete_step_test.go new file mode 100644 index 000000000..50cd3029c --- /dev/null +++ b/internal/config/store_postgres_complete_step_test.go @@ -0,0 +1,278 @@ +package config + +// store_postgres_complete_step_test.go -- pgxmock regression tests for +// CompletePlanStep, the atomic ramp-step advance added for issue #1071 and made +// idempotent per step for issue #1669. +// +// The #1071 bug: updatePlanProgress did a plain GetPurchasePlan -> +// CurrentStep++ -> UpdatePurchasePlan read-modify-write with no row lock, so +// two overlapping Lambda invocations could both read CurrentStep=N and both +// write N+1, skipping a ramp step (a lost update on a money path). +// +// The #1669 bug: even under the lock, the advance was a blind ++ that carried +// no notion of WHICH step completed. A multi-account plan produces one +// execution per account per ramp step, so retrying two separately-failed +// accounts of one step advanced the ramp twice and the plan reported itself a +// step further along than the commitment it had bought. +// +// These tests pin the fix's contract: +// - the read happens inside a transaction (Begin) and carries FOR UPDATE, +// - completing step N from CurrentStep N-1 persists CurrentStep = N, +// - completing a step the ramp already passed writes nothing at all, +// - completing a step more than one beyond CurrentStep is refused, +// - a non-positive step is refused before the transaction is opened, +// - a plan deleted mid-race is tolerated (returns nil, no spurious error), +// - a completed ramp clears next_execution_date instead of advancing. +// +// pgxmock uses regexp query matching (see newMock), so an expectation whose +// pattern requires "FOR UPDATE" only matches if the production query actually +// emits it -- that is the atomicity guard. Dropping FOR UPDATE from the store +// makes TestPGXMock_CompletePlanStep_LocksAndAdvances fail. + +import ( + "context" + "database/sql" + "encoding/json" + "testing" + "time" + + "github.com/jackc/pgx/v5" + "github.com/pashagolub/pgxmock/v4" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// purchasePlanRowCols mirrors the SELECT column list in purchasePlanSelectCols. +var purchasePlanRowCols = []string{ + "id", "name", "enabled", "auto_purchase", "notification_days_before", + "services", "ramp_schedule", "created_at", "updated_at", + "next_execution_date", "last_execution_date", "last_notification_sent", +} + +// rampStepArg is a pgxmock.Argument that unmarshals the ramp_schedule JSONB +// passed to the UPDATE and asserts its CurrentStep, proving the advance landed +// in the persisted row rather than only in memory. +type rampStepArg struct{ want int } + +func (a rampStepArg) Match(v interface{}) bool { + b, ok := v.([]byte) + if !ok { + return false + } + var rs RampSchedule + if err := json.Unmarshal(b, &rs); err != nil { + return false + } + return rs.CurrentStep == a.want +} + +// nullTimeArg matches a *time.Time UPDATE argument by presence (nil vs set), +// used to assert next_execution_date is cleared on a completed ramp. +type nullTimeArg struct{ wantNil bool } + +func (a nullTimeArg) Match(v interface{}) bool { + tp, ok := v.(*time.Time) + if !ok { + return false + } + return (tp == nil) == a.wantNil +} + +// afterArg matches a time.Time UPDATE argument that is strictly after the given +// instant, used to assert updated_at is refreshed rather than persisted stale. +type afterArg struct{ notBefore time.Time } + +func (a afterArg) Match(v interface{}) bool { + ts, ok := v.(time.Time) + if !ok { + return false + } + return ts.After(a.notBefore) +} + +const purchasePlanUpdateArgs = 11 + +// completeStepUpdateArgs builds the WithArgs matcher list for the UPDATE issued +// by CompletePlanStep, asserting the persisted CurrentStep at $7, a refreshed +// updated_at at $8 (> staleUpdatedAt), and the next_execution_date presence at +// $9 while leaving the rest as AnyArg. +func completeStepUpdateArgs(wantStep int, staleUpdatedAt time.Time, wantNextNil bool) []interface{} { + args := anyArgsCfg(purchasePlanUpdateArgs) + args[6] = rampStepArg{want: wantStep} // ramp_schedule = $7 + args[7] = afterArg{notBefore: staleUpdatedAt} // updated_at = $8 + args[8] = nullTimeArg{wantNil: wantNextNil} // next_execution_date = $9 + return args +} + +// rampPlanRows builds a single purchase_plans row carrying ramp, with +// updated_at seeded to stale so a test can assert the store refreshes it. +func rampPlanRows(t *testing.T, planID string, ramp RampSchedule, now, stale time.Time, nextExec sql.NullTime) *pgxmock.Rows { + t.Helper() + svcJSON, err := json.Marshal(map[string]ServiceConfig{}) + require.NoError(t, err) + rampJSON, err := json.Marshal(ramp) + require.NoError(t, err) + + return pgxmock.NewRows(purchasePlanRowCols).AddRow( + planID, "Ramp Plan", true, true, 3, + svcJSON, rampJSON, now, stale, + nextExec, sql.NullTime{Valid: false}, sql.NullTime{Valid: false}, + ) +} + +func TestPGXMock_CompletePlanStep_LocksAndAdvances(t *testing.T) { + mock := newMock(t) + store := storeWith(mock) + ctx := context.Background() + + now := time.Now().Truncate(time.Second) + stale := now.AddDate(0, 0, -30) + ramp := RampSchedule{ + Type: "weekly", + PercentPerStep: 25, + StepIntervalDays: 7, + CurrentStep: 1, + TotalSteps: 4, + StartDate: now, + } + + // Begin -> SELECT ... FOR UPDATE -> UPDATE (CurrentStep advanced 1 -> 2, + // updated_at refreshed, next_execution_date still set since ramp not + // complete) -> Commit. + mock.ExpectBegin() + mock.ExpectQuery(`SELECT[\s\S]*FROM purchase_plans[\s\S]*WHERE id = \$1 FOR UPDATE`). + WithArgs("plan-123"). + WillReturnRows(rampPlanRows(t, "plan-123", ramp, now, stale, sql.NullTime{Valid: false})) + mock.ExpectExec(`UPDATE purchase_plans`). + WithArgs(completeStepUpdateArgs(2, stale, false)...). + WillReturnResult(pgxmock.NewResult("UPDATE", 1)) + mock.ExpectCommit() + + require.NoError(t, store.CompletePlanStep(ctx, "plan-123", 2)) + require.NoError(t, mock.ExpectationsWereMet()) +} + +// TestPGXMock_CompletePlanStep_AlreadyCountedStepWritesNothing is the issue +// #1669 guard at the store boundary: the second per-account retry of one ramp +// step completes a step the plan already counted, and must leave the row +// exactly as it found it. pgxmock fails the test if any UPDATE is issued, +// because no ExpectExec is registered. +func TestPGXMock_CompletePlanStep_AlreadyCountedStepWritesNothing(t *testing.T) { + mock := newMock(t) + store := storeWith(mock) + ctx := context.Background() + + now := time.Now().Truncate(time.Second) + stale := now.AddDate(0, 0, -30) + // The ramp already advanced to 3 when the first retry of step 3 landed. + ramp := RampSchedule{Type: "weekly", PercentPerStep: 25, StepIntervalDays: 7, CurrentStep: 3, TotalSteps: 4, StartDate: now} + + mock.ExpectBegin() + mock.ExpectQuery(`SELECT[\s\S]*FROM purchase_plans[\s\S]*WHERE id = \$1 FOR UPDATE`). + WithArgs("plan-123"). + WillReturnRows(rampPlanRows(t, "plan-123", ramp, now, stale, sql.NullTime{Valid: true, Time: now})) + mock.ExpectCommit() + + require.NoError(t, store.CompletePlanStep(ctx, "plan-123", 3), + "re-completing a counted step is a no-op, not an error") + require.NoError(t, mock.ExpectationsWereMet()) +} + +// TestPGXMock_CompletePlanStep_SkippedPredecessorIsRefused pins the other half +// of the idempotency contract: advancing over steps that never completed would +// overstate how much commitment the plan has bought, so it errors rather than +// jumping. No UPDATE is expected. +func TestPGXMock_CompletePlanStep_SkippedPredecessorIsRefused(t *testing.T) { + mock := newMock(t) + store := storeWith(mock) + ctx := context.Background() + + now := time.Now().Truncate(time.Second) + stale := now.AddDate(0, 0, -30) + ramp := RampSchedule{Type: "weekly", PercentPerStep: 25, StepIntervalDays: 7, CurrentStep: 1, TotalSteps: 6, StartDate: now} + + mock.ExpectBegin() + mock.ExpectQuery(`SELECT[\s\S]*FROM purchase_plans[\s\S]*WHERE id = \$1 FOR UPDATE`). + WithArgs("plan-gap"). + WillReturnRows(rampPlanRows(t, "plan-gap", ramp, now, stale, sql.NullTime{Valid: false})) + mock.ExpectRollback() + + err := store.CompletePlanStep(ctx, "plan-gap", 4) + require.Error(t, err, "steps 2-3 never completed, so step 4 must not silently advance the ramp") + assert.Contains(t, err.Error(), "never completed") + require.NoError(t, mock.ExpectationsWereMet()) +} + +// TestPGXMock_CompletePlanStep_NonPositiveStepIsRefused asserts the argument is +// validated before any transaction opens: a plan-attributed execution that +// carries no ramp step must not advance the ramp by guesswork. +func TestPGXMock_CompletePlanStep_NonPositiveStepIsRefused(t *testing.T) { + mock := newMock(t) + store := storeWith(mock) + ctx := context.Background() + + // No Begin is expected: the guard must short-circuit before the round trip. + err := store.CompletePlanStep(ctx, "plan-123", 0) + require.Error(t, err) + assert.Contains(t, err.Error(), "1-based") + require.NoError(t, mock.ExpectationsWereMet()) +} + +func TestPGXMock_CompletePlanStep_PlanDeletedMidRace(t *testing.T) { + mock := newMock(t) + store := storeWith(mock) + ctx := context.Background() + + mock.ExpectBegin() + mock.ExpectQuery(`SELECT[\s\S]*FROM purchase_plans[\s\S]*WHERE id = \$1 FOR UPDATE`). + WithArgs("gone"). + WillReturnError(pgx.ErrNoRows) + mock.ExpectCommit() + + // A plan deleted between execution and progress update must not error: the + // caller cannot control that race and should not be penalized for it. + require.NoError(t, store.CompletePlanStep(ctx, "gone", 2)) + require.NoError(t, mock.ExpectationsWereMet()) +} + +func TestPGXMock_CompletePlanStep_CompletedRampClearsNextDate(t *testing.T) { + mock := newMock(t) + store := storeWith(mock) + ctx := context.Background() + + now := time.Now().Truncate(time.Second) + stale := now.AddDate(0, 0, -30) + // Already at the last step: CurrentStep == TotalSteps means IsComplete, so + // the step is not advanced and next_execution_date is cleared. + ramp := RampSchedule{Type: "weekly", PercentPerStep: 25, StepIntervalDays: 7, CurrentStep: 4, TotalSteps: 4, StartDate: now} + + mock.ExpectBegin() + mock.ExpectQuery(`SELECT[\s\S]*FROM purchase_plans[\s\S]*WHERE id = \$1 FOR UPDATE`). + WithArgs("plan-done"). + WillReturnRows(rampPlanRows(t, "plan-done", ramp, now, stale, sql.NullTime{Valid: true, Time: now})) + // Step stays at 4 (already complete), updated_at refreshed, and + // next_execution_date is cleared. + mock.ExpectExec(`UPDATE purchase_plans`). + WithArgs(completeStepUpdateArgs(4, stale, true)...). + WillReturnResult(pgxmock.NewResult("UPDATE", 1)) + mock.ExpectCommit() + + require.NoError(t, store.CompletePlanStep(ctx, "plan-done", 5)) + require.NoError(t, mock.ExpectationsWereMet()) +} + +func TestPGXMock_CompletePlanStep_LockErrorRollsBack(t *testing.T) { + mock := newMock(t) + store := storeWith(mock) + ctx := context.Background() + + mock.ExpectBegin() + mock.ExpectQuery(`SELECT[\s\S]*FROM purchase_plans[\s\S]*WHERE id = \$1 FOR UPDATE`). + WithArgs("plan-err"). + WillReturnError(assert.AnError) + mock.ExpectRollback() + + err := store.CompletePlanStep(ctx, "plan-err", 2) + require.Error(t, err) + require.NoError(t, mock.ExpectationsWereMet()) +} diff --git a/internal/config/store_postgres_increment_step_test.go b/internal/config/store_postgres_increment_step_test.go deleted file mode 100644 index ac8897b93..000000000 --- a/internal/config/store_postgres_increment_step_test.go +++ /dev/null @@ -1,211 +0,0 @@ -package config - -// store_postgres_increment_step_test.go -- pgxmock regression tests for -// IncrementPlanCurrentStep, the atomic ramp-step advance added for issue #1071. -// -// The bug: updatePlanProgress previously did a plain GetPurchasePlan -> -// CurrentStep++ -> UpdatePurchasePlan read-modify-write with no row lock, so -// two overlapping Lambda invocations could both read CurrentStep=N and both -// write N+1, skipping a ramp step (a lost update on a money path). -// -// These tests pin the fix's contract: -// - the read happens inside a transaction (Begin) and carries FOR UPDATE, -// - CurrentStep is advanced by exactly one and that value is persisted, -// - a plan deleted mid-race is tolerated (returns nil, no spurious error), -// - a completed ramp clears next_execution_date instead of advancing. -// -// pgxmock uses regexp query matching (see newMock), so an expectation whose -// pattern requires "FOR UPDATE" only matches if the production query actually -// emits it -- that is the atomicity guard. Dropping FOR UPDATE from the store -// makes TestPGXMock_IncrementPlanCurrentStep_LocksAndAdvances fail. - -import ( - "context" - "database/sql" - "encoding/json" - "testing" - "time" - - "github.com/jackc/pgx/v5" - "github.com/pashagolub/pgxmock/v4" - "github.com/stretchr/testify/assert" - "github.com/stretchr/testify/require" -) - -// purchasePlanRowCols mirrors the SELECT column list in purchasePlanSelectCols. -var purchasePlanRowCols = []string{ - "id", "name", "enabled", "auto_purchase", "notification_days_before", - "services", "ramp_schedule", "created_at", "updated_at", - "next_execution_date", "last_execution_date", "last_notification_sent", -} - -// rampStepArg is a pgxmock.Argument that unmarshals the ramp_schedule JSONB -// passed to the UPDATE and asserts its CurrentStep, proving the advance landed -// in the persisted row rather than only in memory. -type rampStepArg struct{ want int } - -func (a rampStepArg) Match(v interface{}) bool { - b, ok := v.([]byte) - if !ok { - return false - } - var rs RampSchedule - if err := json.Unmarshal(b, &rs); err != nil { - return false - } - return rs.CurrentStep == a.want -} - -// nullTimeArg matches a *time.Time UPDATE argument by presence (nil vs set), -// used to assert next_execution_date is cleared on a completed ramp. -type nullTimeArg struct{ wantNil bool } - -func (a nullTimeArg) Match(v interface{}) bool { - tp, ok := v.(*time.Time) - if !ok { - return false - } - return (tp == nil) == a.wantNil -} - -// afterArg matches a time.Time UPDATE argument that is strictly after the given -// instant, used to assert updated_at is refreshed rather than persisted stale. -type afterArg struct{ notBefore time.Time } - -func (a afterArg) Match(v interface{}) bool { - ts, ok := v.(time.Time) - if !ok { - return false - } - return ts.After(a.notBefore) -} - -const purchasePlanUpdateArgs = 11 - -// incrementUpdateArgs builds the WithArgs matcher list for the UPDATE issued by -// IncrementPlanCurrentStep, asserting the persisted CurrentStep at $7, a -// refreshed updated_at at $8 (> staleUpdatedAt), and the next_execution_date -// presence at $9 while leaving the rest as AnyArg. -func incrementUpdateArgs(wantStep int, staleUpdatedAt time.Time, wantNextNil bool) []interface{} { - args := anyArgsCfg(purchasePlanUpdateArgs) - args[6] = rampStepArg{want: wantStep} // ramp_schedule = $7 - args[7] = afterArg{notBefore: staleUpdatedAt} // updated_at = $8 - args[8] = nullTimeArg{wantNil: wantNextNil} // next_execution_date = $9 - return args -} - -func TestPGXMock_IncrementPlanCurrentStep_LocksAndAdvances(t *testing.T) { - mock := newMock(t) - store := storeWith(mock) - ctx := context.Background() - - now := time.Now().Truncate(time.Second) - svcJSON, err := json.Marshal(map[string]ServiceConfig{}) - require.NoError(t, err) - ramp := RampSchedule{ - Type: "weekly", - PercentPerStep: 25, - StepIntervalDays: 7, - CurrentStep: 1, - TotalSteps: 4, - StartDate: now, - } - rampJSON, err := json.Marshal(ramp) - require.NoError(t, err) - - // Seed updated_at with an old timestamp so the test can assert the store - // refreshes it on increment rather than persisting the stale value. - stale := now.AddDate(0, 0, -30) - rows := pgxmock.NewRows(purchasePlanRowCols).AddRow( - "plan-123", "Ramp Plan", true, true, 3, - svcJSON, rampJSON, now, stale, - sql.NullTime{Valid: false}, sql.NullTime{Valid: false}, sql.NullTime{Valid: false}, - ) - - // Begin -> SELECT ... FOR UPDATE -> UPDATE (CurrentStep advanced 1 -> 2, - // updated_at refreshed, next_execution_date still set since ramp not - // complete) -> Commit. - mock.ExpectBegin() - mock.ExpectQuery(`SELECT[\s\S]*FROM purchase_plans[\s\S]*WHERE id = \$1 FOR UPDATE`). - WithArgs("plan-123"). - WillReturnRows(rows) - mock.ExpectExec(`UPDATE purchase_plans`). - WithArgs(incrementUpdateArgs(2, stale, false)...). - WillReturnResult(pgxmock.NewResult("UPDATE", 1)) - mock.ExpectCommit() - - err = store.IncrementPlanCurrentStep(ctx, "plan-123") - require.NoError(t, err) - require.NoError(t, mock.ExpectationsWereMet()) -} - -func TestPGXMock_IncrementPlanCurrentStep_PlanDeletedMidRace(t *testing.T) { - mock := newMock(t) - store := storeWith(mock) - ctx := context.Background() - - mock.ExpectBegin() - mock.ExpectQuery(`SELECT[\s\S]*FROM purchase_plans[\s\S]*WHERE id = \$1 FOR UPDATE`). - WithArgs("gone"). - WillReturnError(pgx.ErrNoRows) - mock.ExpectCommit() - - // A plan deleted between execution and progress update must not error: the - // caller cannot control that race and should not be penalized for it. - err := store.IncrementPlanCurrentStep(ctx, "gone") - require.NoError(t, err) - require.NoError(t, mock.ExpectationsWereMet()) -} - -func TestPGXMock_IncrementPlanCurrentStep_CompletedRampClearsNextDate(t *testing.T) { - mock := newMock(t) - store := storeWith(mock) - ctx := context.Background() - - now := time.Now().Truncate(time.Second) - svcJSON, err := json.Marshal(map[string]ServiceConfig{}) - require.NoError(t, err) - // Already at the last step: CurrentStep == TotalSteps means IsComplete, so - // the step is not advanced and next_execution_date is cleared. - ramp := RampSchedule{Type: "weekly", PercentPerStep: 25, StepIntervalDays: 7, CurrentStep: 4, TotalSteps: 4, StartDate: now} - rampJSON, err := json.Marshal(ramp) - require.NoError(t, err) - - stale := now.AddDate(0, 0, -30) - rows := pgxmock.NewRows(purchasePlanRowCols).AddRow( - "plan-done", "Done Plan", true, true, 3, - svcJSON, rampJSON, now, stale, - sql.NullTime{Valid: true, Time: now}, sql.NullTime{Valid: false}, sql.NullTime{Valid: false}, - ) - - mock.ExpectBegin() - mock.ExpectQuery(`SELECT[\s\S]*FROM purchase_plans[\s\S]*WHERE id = \$1 FOR UPDATE`). - WithArgs("plan-done"). - WillReturnRows(rows) - // Step stays at 4 (already complete), updated_at refreshed, and - // next_execution_date is cleared. - mock.ExpectExec(`UPDATE purchase_plans`). - WithArgs(incrementUpdateArgs(4, stale, true)...). - WillReturnResult(pgxmock.NewResult("UPDATE", 1)) - mock.ExpectCommit() - - err = store.IncrementPlanCurrentStep(ctx, "plan-done") - require.NoError(t, err) - require.NoError(t, mock.ExpectationsWereMet()) -} - -func TestPGXMock_IncrementPlanCurrentStep_LockErrorRollsBack(t *testing.T) { - mock := newMock(t) - store := storeWith(mock) - ctx := context.Background() - - mock.ExpectBegin() - mock.ExpectQuery(`SELECT[\s\S]*FROM purchase_plans[\s\S]*WHERE id = \$1 FOR UPDATE`). - WithArgs("plan-err"). - WillReturnError(assert.AnError) - mock.ExpectRollback() - - err := store.IncrementPlanCurrentStep(ctx, "plan-err") - require.Error(t, err) - require.NoError(t, mock.ExpectationsWereMet()) -} diff --git a/internal/database/postgres/migrations/000098_backfill_execution_step_number.down.sql b/internal/database/postgres/migrations/000098_backfill_execution_step_number.down.sql new file mode 100644 index 000000000..b16018f6e --- /dev/null +++ b/internal/database/postgres/migrations/000098_backfill_execution_step_number.down.sql @@ -0,0 +1,13 @@ +-- 000098 down: intentional no-op. +-- +-- The up migration rewrites purchase_executions.step_number for executable +-- rows that carried the pre-#1669 convention. The rows it changed cannot be +-- identified afterwards, so an inverse UPDATE would also rewrite correctly +-- stamped rows and corrupt them. +-- +-- Leaving the corrected values in place is safe for a rollback: the pre-#1669 +-- code advances the ramp with a blind CurrentStep++ and never reads +-- step_number to decide anything, so a corrected row behaves identically under +-- the old code. A no-op makes that explicit rather than hiding it behind a +-- silent empty file. +SELECT 1; -- no-op: the affected rows cannot be identified after the fact diff --git a/internal/database/postgres/migrations/000098_backfill_execution_step_number.up.sql b/internal/database/postgres/migrations/000098_backfill_execution_step_number.up.sql new file mode 100644 index 000000000..59119b8cd --- /dev/null +++ b/internal/database/postgres/migrations/000098_backfill_execution_step_number.up.sql @@ -0,0 +1,65 @@ +-- Migration 000098: backfill purchase_executions.step_number onto the 1-based +-- "step this row completes" convention (issue #1669). +-- +-- WHY +-- purchase.getOrCreateExecution stamped step_number with the COUNT of ramp +-- steps already completed (plan.ramp_schedule.current_step), while +-- api.createPurchaseExecutionsTx stamped the 1-based step the row would +-- execute (current_step + i + 1). The disagreement was invisible while the +-- ramp advanced by a blind CurrentStep++ that ignored step_number entirely. +-- +-- Issue #1669 makes the advance idempotent per step: it advances only when the +-- completing step is exactly current_step + 1. An old-convention row therefore +-- completes a step the plan has already counted, which the new code correctly +-- treats as a no-op, so the ramp would never advance again for that plan and +-- the remaining steps would silently never purchase. That is the same +-- under-buy shape #1669 exists to fix, so the in-flight rows must be corrected +-- as part of it. +-- +-- WHICH ROWS +-- Only rows that can still execute. A correctly-stamped executable row names a +-- step the plan has not yet counted, i.e. step_number >= current_step + 1, so +-- step_number <= current_step is the old-convention marker. +-- +-- One correctly-stamped shape defeats that marker on its own, hence the NOT +-- EXISTS guard. A multi-account step that partially failed leaves per-account +-- retry successors carrying the step they belong to; once the first retry +-- advances the ramp, a sibling successor still pending for the SAME step now +-- reads step_number = current_step. Retargeting it would make it buy its own +-- step's tranche while advancing the ramp past the next one, which is exactly +-- the under-buy #1669 fixes. Such a row is recognizable: a sibling execution +-- for that plan and step has already reached a successful terminal status, +-- which is never true of an old-convention row (the old writer stamped the row +-- that advanced step k as k-1, so no completed row carries k). +-- +-- 'paused' counts as executable: resumePlannedPurchase transitions it back to +-- 'pending' and runPlannedPurchase transitions it straight to 'running' +-- (internal/api/handler_purchases.go). +-- +-- Terminal rows (completed / partially_completed / failed / canceled / +-- expired / revocation_requested) are deliberately NOT rewritten: their +-- step_number is audit trail and display, and no discriminator distinguishes +-- an old-convention terminal row from a correctly-stamped one whose ramp has +-- since moved past it. +-- +-- DEPLOY WINDOW +-- Old code instances still running when this migration lands keep writing the +-- old convention, exactly as migration 000089 documents for the cancelled -> +-- canceled rename. Such a row freezes its own plan's ramp at one step until an +-- operator intervenes. The window is the deploy overlap, and the alternative +-- (deriving current_step from the executions table) is a larger change tracked +-- separately. + +UPDATE purchase_executions e + SET step_number = COALESCE((p.ramp_schedule ->> 'current_step')::int, 0) + 1 + FROM purchase_plans p + WHERE e.plan_id = p.id + AND e.status IN ('pending', 'notified', 'approved', 'running', 'scheduled', 'paused') + AND e.step_number <= COALESCE((p.ramp_schedule ->> 'current_step')::int, 0) + AND NOT EXISTS ( + SELECT 1 + FROM purchase_executions sibling + WHERE sibling.plan_id = e.plan_id + AND sibling.step_number = e.step_number + AND sibling.status IN ('completed', 'partially_completed') + ); diff --git a/internal/database/postgres/migrations/000098_backfill_execution_step_number_test.go b/internal/database/postgres/migrations/000098_backfill_execution_step_number_test.go new file mode 100644 index 000000000..ed2f6e24d --- /dev/null +++ b/internal/database/postgres/migrations/000098_backfill_execution_step_number_test.go @@ -0,0 +1,170 @@ +//go:build integration +// +build integration + +package migrations_test + +import ( + "context" + "testing" + + "github.com/LeanerCloud/CUDly/internal/database/postgres/migrations" + "github.com/LeanerCloud/CUDly/internal/database/postgres/testhelpers" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +// seedStepNumberPlan inserts a purchase_plans row sitting on ramp step +// currentStep and returns its id. +func seedStepNumberPlan(ctx context.Context, t *testing.T, pool *pgxpool.Pool, name string, currentStep int) string { + t.Helper() + var planID string + err := pool.QueryRow(ctx, ` + INSERT INTO purchase_plans (name, ramp_schedule) + VALUES ($1, jsonb_build_object('type', 'weekly', 'percent_per_step', 25, + 'step_interval_days', 7, 'total_steps', 4, + 'current_step', $2::int)) + RETURNING id + `, name, currentStep).Scan(&planID) + require.NoError(t, err) + return planID +} + +// seedStepNumberExecution inserts a purchase_executions row and returns its id. +func seedStepNumberExecution(ctx context.Context, t *testing.T, pool *pgxpool.Pool, planID, status string, stepNumber int) string { + t.Helper() + var execID string + err := pool.QueryRow(ctx, ` + INSERT INTO purchase_executions (plan_id, status, step_number, scheduled_date) + VALUES ($1, $2, $3, NOW()) + RETURNING execution_id + `, planID, status, stepNumber).Scan(&execID) + require.NoError(t, err) + return execID +} + +// stepNumberOf reads back one execution's step_number. +func stepNumberOf(ctx context.Context, t *testing.T, pool *pgxpool.Pool, execID string) int { + t.Helper() + var step int + require.NoError(t, pool.QueryRow(ctx, + `SELECT step_number FROM purchase_executions WHERE execution_id = $1`, execID).Scan(&step)) + return step +} + +// TestMigration_BackfillExecutionStepNumber locks down migration 000098: rows +// carrying the pre-#1669 convention (step_number = the COUNT of completed ramp +// steps) are retargeted at the step they will actually complete, while +// correctly-stamped rows and terminal rows are left alone. +func TestMigration_BackfillExecutionStepNumber(t *testing.T) { + ctx := context.Background() + migrationsPath := getMigrationsPath() + + container, err := testhelpers.SetupPostgresContainer(ctx, t) + require.NoError(t, err) + defer container.Cleanup(ctx) + pool := container.DB.Pool() + + // Pin just below 000098 so the assertions exercise this migration's direct + // effect. The number below head is read from disk because migration numbers + // in this repo are not contiguous. + require.NoError(t, migrations.MigrateToVersion(ctx, pool, migrationsPath, previousMigrationVersion(t, 98))) + + // A fresh plan (nothing bought yet) whose pending row was stamped 0 by the + // old writer: the shape that would hit the "no ramp step" refusal. + freshPlan := seedStepNumberPlan(ctx, t, pool, "Fresh Plan", 0) + freshPending := seedStepNumberExecution(ctx, t, pool, freshPlan, "pending", 0) + + // A plan two steps in whose approved row was stamped 2 by the old writer: + // the shape that would silently no-op forever. + midPlan := seedStepNumberPlan(ctx, t, pool, "Mid Plan", 2) + midApproved := seedStepNumberExecution(ctx, t, pool, midPlan, "approved", 2) + // Correctly stamped rows on the same plan must not move. + midCorrect := seedStepNumberExecution(ctx, t, pool, midPlan, "pending", 3) + midFuture := seedStepNumberExecution(ctx, t, pool, midPlan, "notified", 4) + // The old-convention row that advanced this plan 1 -> 2, stamped 1 by the + // old writer. It must not be rewritten, but note WHICH clause spares it: a + // completed row is its own completed sibling, so NOT EXISTS excludes it + // whether or not the status allowlist is there. Its job in this fixture is + // to establish that no completed row carries step 2, which is what lets + // midApproved and midFailed below isolate the other two clauses. + midTerminal := seedStepNumberExecution(ctx, t, pool, midPlan, "completed", 1) + // The row that isolates the status allowlist. It is terminal but NOT + // completed, so NOT EXISTS is true for it (no completed row carries step 2) + // and step_number <= current_step holds: line 57 of the migration is the + // only thing excluding it. Rewriting it would not be cosmetic, because + // persistRetryExecution propagates a failed row's StepNumber onto its retry + // successor, so moving it changes which ramp step a later operator retry + // completes. + midFailed := seedStepNumberExecution(ctx, t, pool, midPlan, "failed", 2) + // Every other status the plan flow can still execute from. Missing one is + // how 'paused' survived three review passes, so cover the set explicitly. + midExecutable := map[string]string{} + for _, status := range []string{"paused", "running", "scheduled", "notified"} { + midExecutable[status] = seedStepNumberExecution(ctx, t, pool, midPlan, status, 2) + } + + // A plan whose ramp_schedule carries no current_step key at all: the + // COALESCE must read it as 0, so its step-0 row is still retargeted. Without + // the COALESCE the comparison is NULL and the row is silently skipped. + var keylessPlan string + require.NoError(t, pool.QueryRow(ctx, ` + INSERT INTO purchase_plans (name, ramp_schedule) + VALUES ('Keyless Plan', jsonb_build_object('type', 'weekly', 'total_steps', 4)) + RETURNING id + `).Scan(&keylessPlan)) + keylessPending := seedStepNumberExecution(ctx, t, pool, keylessPlan, "pending", 0) + + // The shape the NOT EXISTS guard protects: a multi-account step 3 that + // partially failed, where the first per-account retry already advanced the + // ramp to 3 and a sibling successor for the SAME step is still pending. Its + // step_number is correct; retargeting it to 4 would make it buy step 3's + // tranche while advancing the ramp past step 4. + partialPlan := seedStepNumberPlan(ctx, t, pool, "Partial Step Plan", 3) + seedStepNumberExecution(ctx, t, pool, partialPlan, "completed", 3) + partialSibling := seedStepNumberExecution(ctx, t, pool, partialPlan, "pending", 3) + + // Same shape, but the sibling that advanced the ramp landed + // partially_completed rather than completed. That is the LIKELIER form of + // this scenario, not a corner: applyAccountOutcome stamps + // partially_completed when some of an account's recs committed and others + // failed, which is what a partial fan-out produces. It isolates the second + // entry of the sibling list, which nothing else in this fixture exercises. + partialRecPlan := seedStepNumberPlan(ctx, t, pool, "Partially Completed Sibling Plan", 3) + seedStepNumberExecution(ctx, t, pool, partialRecPlan, "partially_completed", 3) + partialRecSibling := seedStepNumberExecution(ctx, t, pool, partialRecPlan, "pending", 3) + + require.Equal(t, 0, stepNumberOf(ctx, t, pool, freshPending), "seed must start on the old convention") + require.Equal(t, 2, stepNumberOf(ctx, t, pool, midApproved), "seed must start on the old convention") + + require.NoError(t, migrations.MigrateToVersion(ctx, pool, migrationsPath, 98)) + + assert.Equal(t, 1, stepNumberOf(ctx, t, pool, freshPending), + "a step-0 executable row on a plan at step 0 must be retargeted at step 1") + assert.Equal(t, 3, stepNumberOf(ctx, t, pool, midApproved), + "an old-convention executable row must be retargeted at current_step + 1") + assert.Equal(t, 3, stepNumberOf(ctx, t, pool, midCorrect), + "a correctly-stamped row for the next step must not move") + assert.Equal(t, 4, stepNumberOf(ctx, t, pool, midFuture), + "a correctly-stamped row for a later step must not move") + assert.Equal(t, 1, stepNumberOf(ctx, t, pool, midTerminal), + "a completed row is audit trail and must not be rewritten (excluded by NOT EXISTS, which finds itself)") + assert.Equal(t, 2, stepNumberOf(ctx, t, pool, midFailed), + "a failed row is audit trail and a retry successor's step source; only the status allowlist excludes it") + for status, execID := range midExecutable { + assert.Equal(t, 3, stepNumberOf(ctx, t, pool, execID), + "a %q row can still execute and must be retargeted", status) + } + assert.Equal(t, 1, stepNumberOf(ctx, t, pool, keylessPending), + "a plan with no current_step key must be read as step 0, not skipped") + assert.Equal(t, 3, stepNumberOf(ctx, t, pool, partialSibling), + "a pending sibling of an already-completed step is correctly stamped and must not move") + assert.Equal(t, 3, stepNumberOf(ctx, t, pool, partialRecSibling), + "a pending sibling of a partially_completed step must not move either; only that entry of the sibling list excludes it") + + // The down migration is a documented no-op: rolling back the version must + // succeed and must leave the corrected values in place. + require.NoError(t, migrations.MigrateToVersion(ctx, pool, migrationsPath, previousMigrationVersion(t, 98))) + assert.Equal(t, 3, stepNumberOf(ctx, t, pool, midApproved), + "the no-op down must leave the corrected value, which the pre-#1669 code ignores") +} diff --git a/internal/mocks/stores.go b/internal/mocks/stores.go index 2129e9364..6df67fc12 100644 --- a/internal/mocks/stores.go +++ b/internal/mocks/stores.go @@ -191,10 +191,10 @@ func (m *MockConfigStore) UpdatePurchasePlan(ctx context.Context, plan *config.P return args.Error(0) } -// IncrementPlanCurrentStep mocks the atomic step-advance operation. -func (m *MockConfigStore) IncrementPlanCurrentStep(ctx context.Context, planID string) error { - m.record("IncrementPlanCurrentStep", ctx, planID) - args := m.Called(ctx, planID) +// CompletePlanStep mocks the idempotent ramp-step advance operation. +func (m *MockConfigStore) CompletePlanStep(ctx context.Context, planID string, stepNumber int) error { + m.record("CompletePlanStep", ctx, planID, stepNumber) + args := m.Called(ctx, planID, stepNumber) return args.Error(0) } diff --git a/internal/purchase/approvals_test.go b/internal/purchase/approvals_test.go index 54723a8fc..98feb62e4 100644 --- a/internal/purchase/approvals_test.go +++ b/internal/purchase/approvals_test.go @@ -46,8 +46,9 @@ func stubExecuteChain(t *testing.T, store *MockConfigStore, sender *MockEmailSen sender.On("SendPurchaseConfirmation", mock.Anything, mock.Anything).Return(nil) // finalizeExecution writes status=completed via SavePurchaseExecution. store.On("SavePurchaseExecution", mock.Anything, mock.AnythingOfType("*config.PurchaseExecution")).Return(nil) - // updatePlanProgress calls IncrementPlanCurrentStep (atomic, issue #1071). - store.On("IncrementPlanCurrentStep", mock.Anything, planID).Return(nil) + // updatePlanProgress calls CompletePlanStep (atomic, issue #1071); every + // caller of this helper stamps StepNumber: 1 on its "updated" fixture. + store.On("CompletePlanStep", mock.Anything, planID, 1).Return(nil) } func TestManager_ApproveExecution_Success(t *testing.T) { @@ -65,6 +66,7 @@ func TestManager_ApproveExecution_Success(t *testing.T) { PlanID: "plan-456", Status: "approved", ApprovalToken: "valid-token", + StepNumber: 1, } store.On("GetExecutionByID", ctx, "exec-123").Return(execution, nil) @@ -92,6 +94,7 @@ func TestManager_ApproveExecution_StampsApprovedBy(t *testing.T) { PlanID: "plan-456", Status: "approved", ApprovalToken: "valid-token", + StepNumber: 1, } store.On("GetExecutionByID", ctx, "exec-123").Return(execution, nil) @@ -125,6 +128,7 @@ func TestManager_ApproveExecution_NotifiedStatus(t *testing.T) { PlanID: "plan-456", Status: "approved", ApprovalToken: "valid-token", + StepNumber: 1, } store.On("GetExecutionByID", ctx, "exec-123").Return(execution, nil) @@ -288,6 +292,7 @@ func TestManager_ApproveAndExecute_SkipsTokenCheck(t *testing.T) { PlanID: "plan-789", Status: "approved", ApprovalToken: "tok", + StepNumber: 1, } store.On("TransitionExecutionStatus", ctx, "exec-456", approveFromStatuses, "approved", mock.MatchedBy(func(actor *string) bool { return actor != nil && *actor == actorUUID })).Return(updated, nil) @@ -376,6 +381,7 @@ func TestManager_ApproveAndExecute_FourEyesOn_AllowsDifferentApprover(t *testing ExecutionID: "exec-direct-diff", PlanID: "plan-fourEyes", Status: "approved", + StepNumber: 1, } store.On("GetGlobalConfig", ctx).Return(fourEyesCfgOnForManager(), nil) @@ -471,7 +477,7 @@ func TestManager_ApproveAndExecute_FourEyesOff_AllowsSelfApprove(t *testing.T) { manager, store, sender := newApproveManager(t) creatorID := "user-creator" - updated := &config.PurchaseExecution{ExecutionID: "exec-mode-off", PlanID: "plan-fourEyes", Status: "approved"} + updated := &config.PurchaseExecution{ExecutionID: "exec-mode-off", PlanID: "plan-fourEyes", Status: "approved", StepNumber: 1} store.On("TransitionExecutionStatus", ctx, "exec-mode-off", approveFromStatuses, "approved", &creatorID).Return(updated, nil) stubExecuteChain(t, store, sender, "plan-fourEyes") @@ -536,6 +542,7 @@ func TestManager_ApproveExecution_FourEyesOn_DifferentActor_Allowed(t *testing.T PlanID: "plan-fourEyes", Status: "approved", ApprovalToken: "valid-token", + StepNumber: 1, } store.On("GetExecutionByID", ctx, "exec-sqs-diff").Return(execution, nil) @@ -817,6 +824,7 @@ func TestManager_ApproveExecution_ValidTokenWithinTTL(t *testing.T) { PlanID: "plan-live", Status: "approved", ApprovalToken: "valid-token", + StepNumber: 1, } store.On("GetExecutionByID", ctx, "exec-live").Return(execution, nil) store.On("TransitionExecutionStatus", ctx, "exec-live", approveFromStatuses, "approved", (*string)(nil)).Return(updated, nil) @@ -846,6 +854,7 @@ func TestManager_ApproveExecution_NilExpiresAt_LegacyRow(t *testing.T) { ExecutionID: "exec-legacy", PlanID: "plan-legacy", Status: "approved", + StepNumber: 1, } store.On("GetExecutionByID", ctx, "exec-legacy").Return(execution, nil) store.On("TransitionExecutionStatus", ctx, "exec-legacy", approveFromStatuses, "approved", (*string)(nil)).Return(updated, nil) @@ -906,6 +915,7 @@ func TestManager_ApproveExecution_AWSOrphanFallsThrough(t *testing.T) { PlanID: "plan-aws", Status: "approved", ApprovalToken: "valid-token", + StepNumber: 1, } store.On("GetExecutionByID", ctx, "exec-aws-ambient").Return(execution, nil) store.On("TransitionExecutionStatus", ctx, "exec-aws-ambient", approveFromStatuses, "approved", (*string)(nil)).Return(updated, nil) @@ -1018,6 +1028,7 @@ func TestApproveExecution_MintsRevocationToken(t *testing.T) { PlanID: "plan-rotate", Status: "approved", ApprovalToken: "pre-rotate-token", + StepNumber: 1, } // GetExecutionByID is called twice: once in ApproveExecution itself and diff --git a/internal/purchase/armed_redrive_test.go b/internal/purchase/armed_redrive_test.go index 1b8485fa0..10baaab87 100644 --- a/internal/purchase/armed_redrive_test.go +++ b/internal/purchase/armed_redrive_test.go @@ -105,7 +105,7 @@ func armedHarness(t *testing.T, serviceType common.ServiceType) (*Manager, *Mock store.GetPlanAccountsFn = func(_ context.Context, _ string) ([]config.CloudAccount, error) { return nil, nil } store.On("SavePurchaseHistory", mock.Anything, mock.AnythingOfType("*config.PurchaseHistoryRecord")).Return(nil).Maybe() - store.On("IncrementPlanCurrentStep", mock.Anything, mock.Anything).Return(nil).Maybe() + store.On("CompletePlanStep", mock.Anything, mock.Anything, mock.Anything).Return(nil).Maybe() store.On("GetGlobalConfig", mock.Anything).Return(&config.GlobalConfig{}, nil).Maybe() email.On("SendPurchaseConfirmation", mock.Anything, mock.AnythingOfType("email.NotificationData")).Return(nil).Maybe() diff --git a/internal/purchase/coverage_extra_test.go b/internal/purchase/coverage_extra_test.go index d4beba693..e2c96d149 100644 --- a/internal/purchase/coverage_extra_test.go +++ b/internal/purchase/coverage_extra_test.go @@ -200,6 +200,7 @@ func TestHandleExecutePurchase_ApprovedStatus(t *testing.T) { ExecutionID: "exec-approved", PlanID: "plan-approved", Status: "approved", + StepNumber: 1, Recommendations: []config.RecommendationRecord{}, } @@ -212,7 +213,7 @@ func TestHandleExecutePurchase_ApprovedStatus(t *testing.T) { mockStore.On("GetPurchasePlan", ctx, "plan-approved").Return(plan, nil) mockEmail.On("SendPurchaseConfirmation", ctx, mock.AnythingOfType("email.NotificationData")).Return(nil) mockStore.On("SavePurchaseExecution", ctx, mock.AnythingOfType("*config.PurchaseExecution")).Return(nil) - mockStore.On("IncrementPlanCurrentStep", ctx, "plan-approved").Return(nil) + mockStore.On("CompletePlanStep", ctx, "plan-approved", 1).Return(nil) mockSTS.On("GetCallerIdentity", ctx, mock.Anything).Return(nil, errors.New("sts error")) manager := &Manager{ @@ -351,6 +352,7 @@ func TestProcessMessage_ApproveHappyPath(t *testing.T) { PlanID: planID, Status: "approved", ApprovalToken: "correct-token", + StepNumber: 1, Recommendations: exec.Recommendations, } account := &config.CloudAccount{ID: accountID, ContactEmail: "owner@example.com"} @@ -369,7 +371,7 @@ func TestProcessMessage_ApproveHappyPath(t *testing.T) { mockStore.On("GetPurchasePlan", ctx, planID).Return(plan, nil) mockEmail.On("SendPurchaseConfirmation", ctx, mock.Anything).Return(nil) mockStore.On("SavePurchaseExecution", ctx, mock.AnythingOfType("*config.PurchaseExecution")).Return(nil) - mockStore.On("IncrementPlanCurrentStep", ctx, planID).Return(nil) + mockStore.On("CompletePlanStep", ctx, planID, 1).Return(nil) manager := &Manager{ config: mockStore, @@ -460,6 +462,7 @@ func TestProcessMessage_ApproveFourEyesOn_DifferentApproverSucceeds(t *testing.T ExecutionID: "exec-appv-diff", PlanID: planID, Status: "approved", + StepNumber: 1, Recommendations: exec.Recommendations, } account := &config.CloudAccount{ID: accountID, ContactEmail: "approver@example.com"} @@ -476,7 +479,7 @@ func TestProcessMessage_ApproveFourEyesOn_DifferentApproverSucceeds(t *testing.T mockStore.On("GetPurchasePlan", ctx, planID).Return(plan, nil) mockEmail.On("SendPurchaseConfirmation", ctx, mock.Anything).Return(nil) mockStore.On("SavePurchaseExecution", ctx, mock.AnythingOfType("*config.PurchaseExecution")).Return(nil) - mockStore.On("IncrementPlanCurrentStep", ctx, planID).Return(nil) + mockStore.On("CompletePlanStep", ctx, planID, 1).Return(nil) manager := &Manager{ config: mockStore, diff --git a/internal/purchase/execution.go b/internal/purchase/execution.go index b9dcd9b0a..32210f746 100644 --- a/internal/purchase/execution.go +++ b/internal/purchase/execution.go @@ -1205,14 +1205,28 @@ func mapSavingsPlansSlug(service string) (common.ServiceType, bool) { // Direct-execute purchases (Opportunities flow) have no plan to advance -- // PlanID is empty and the Postgres UUID column would reject the query // with SQLSTATE 22P02, so short-circuit cleanly before hitting the store. -// The actual increment is delegated to IncrementPlanCurrentStep which holds -// a SELECT FOR UPDATE lock for the duration, preventing the lost-update race -// when two Lambda invocations overlap on the same plan (issue #1071). -func (m *Manager) updatePlanProgress(ctx context.Context, planID string) error { - if planID == "" { +// +// The step being completed is read off the execution itself, never off the +// plan's current position: a multi-account step fans out into one execution +// per account, each retryable on its own, so the plan's position is exactly +// what a second completion of one step must not be trusted to imply (issue +// #1669). CompletePlanStep holds a SELECT FOR UPDATE lock for the duration, +// which also closes the lost-update race between overlapping Lambda +// invocations on the same plan (issue #1071). +func (m *Manager) updatePlanProgress(ctx context.Context, exec *config.PurchaseExecution) error { + if exec.PlanID == "" { return nil } - return m.config.IncrementPlanCurrentStep(ctx, planID) + if exec.StepNumber <= 0 { + // Reachable for a row stamped by pre-#1669 code, which wrote the COUNT + // of completed steps and so wrote 0 for a plan's first step. Migration + // 000098 retargets those rows, leaving only ones written during the + // deploy overlap. Advancing the ramp from an unknown step would be a + // guess on the money path; report instead. + return fmt.Errorf("execution %s is attributed to plan %s but carries ramp step %d: refusing to advance the ramp", + exec.ExecutionID, exec.PlanID, exec.StepNumber) + } + return m.config.CompletePlanStep(ctx, exec.PlanID, exec.StepNumber) } // getAWSAccountID retrieves the current AWS account ID using STS. diff --git a/internal/purchase/execution_test.go b/internal/purchase/execution_test.go index 95add7547..67e79549a 100644 --- a/internal/purchase/execution_test.go +++ b/internal/purchase/execution_test.go @@ -322,7 +322,7 @@ func TestManager_UpdatePlanProgress(t *testing.T) { mockStore := new(MockConfigStore) mockEmail := new(MockEmailSender) - mockStore.On("IncrementPlanCurrentStep", ctx, "plan-123").Return(nil) + mockStore.On("CompletePlanStep", ctx, "plan-123", 2).Return(nil) manager := &Manager{ config: mockStore, @@ -330,7 +330,11 @@ func TestManager_UpdatePlanProgress(t *testing.T) { dashboardURL: "https://dashboard.example.com", } - err := manager.updatePlanProgress(ctx, "plan-123") + err := manager.updatePlanProgress(ctx, &config.PurchaseExecution{ + ExecutionID: "exec-123", + PlanID: "plan-123", + StepNumber: 2, + }) require.NoError(t, err) mockStore.AssertExpectations(t) @@ -341,8 +345,8 @@ func TestManager_UpdatePlanProgress_PlanNotFound(t *testing.T) { mockStore := new(MockConfigStore) mockEmail := new(MockEmailSender) - // IncrementPlanCurrentStep returns nil when the plan no longer exists. - mockStore.On("IncrementPlanCurrentStep", ctx, "nonexistent").Return(nil) + // CompletePlanStep returns nil when the plan no longer exists. + mockStore.On("CompletePlanStep", ctx, "nonexistent", 1).Return(nil) manager := &Manager{ config: mockStore, @@ -350,7 +354,11 @@ func TestManager_UpdatePlanProgress_PlanNotFound(t *testing.T) { dashboardURL: "https://dashboard.example.com", } - err := manager.updatePlanProgress(ctx, "nonexistent") + err := manager.updatePlanProgress(ctx, &config.PurchaseExecution{ + ExecutionID: "exec-nonexistent", + PlanID: "nonexistent", + StepNumber: 1, + }) require.NoError(t, err) mockStore.AssertExpectations(t) @@ -361,7 +369,7 @@ func TestManager_UpdatePlanProgress_GetError(t *testing.T) { mockStore := new(MockConfigStore) mockEmail := new(MockEmailSender) - mockStore.On("IncrementPlanCurrentStep", ctx, "plan-123").Return(errors.New("database error")) + mockStore.On("CompletePlanStep", ctx, "plan-123", 1).Return(errors.New("database error")) manager := &Manager{ config: mockStore, @@ -369,7 +377,11 @@ func TestManager_UpdatePlanProgress_GetError(t *testing.T) { dashboardURL: "https://dashboard.example.com", } - err := manager.updatePlanProgress(ctx, "plan-123") + err := manager.updatePlanProgress(ctx, &config.PurchaseExecution{ + ExecutionID: "exec-123", + PlanID: "plan-123", + StepNumber: 1, + }) assert.Error(t, err) mockStore.AssertExpectations(t) @@ -380,9 +392,9 @@ func TestManager_UpdatePlanProgress_CompleteRamp(t *testing.T) { mockStore := new(MockConfigStore) mockEmail := new(MockEmailSender) - // The step-advance and completion logic now live inside IncrementPlanCurrentStep + // The step-completion and ramp-advance logic now live inside CompletePlanStep // in the store, tested at the store layer with pgxmock. - mockStore.On("IncrementPlanCurrentStep", ctx, "plan-123").Return(nil) + mockStore.On("CompletePlanStep", ctx, "plan-123", 4).Return(nil) manager := &Manager{ config: mockStore, @@ -390,12 +402,48 @@ func TestManager_UpdatePlanProgress_CompleteRamp(t *testing.T) { dashboardURL: "https://dashboard.example.com", } - err := manager.updatePlanProgress(ctx, "plan-123") + err := manager.updatePlanProgress(ctx, &config.PurchaseExecution{ + ExecutionID: "exec-123", + PlanID: "plan-123", + StepNumber: 4, + }) require.NoError(t, err) mockStore.AssertExpectations(t) } +// TestManager_UpdatePlanProgress_ZeroStepRefusesAndSkipsStore is the +// regression guard for issue #1669: an execution attributed to a plan but +// carrying a non-positive ramp step must not guess at the ramp position by +// calling the store. updatePlanProgress must return an error and the store +// must never be called. +func TestManager_UpdatePlanProgress_ZeroStepRefusesAndSkipsStore(t *testing.T) { + ctx := context.Background() + mockStore := new(MockConfigStore) + mockEmail := new(MockEmailSender) + + // Registered so AssertNotCalled below is a meaningful check rather than a + // vacuous one (a bare AssertNotCalled without a matching registered + // expectation always passes in this repo's mock setup). + mockStore.On("CompletePlanStep", mock.Anything, mock.Anything, mock.Anything).Return(nil).Maybe() + t.Cleanup(func() { mockStore.AssertExpectations(t) }) + + manager := &Manager{ + config: mockStore, + email: mockEmail, + dashboardURL: "https://dashboard.example.com", + } + + err := manager.updatePlanProgress(ctx, &config.PurchaseExecution{ + ExecutionID: "exec-zero-step", + PlanID: "plan-123", + StepNumber: 0, + }) + require.Error(t, err) + + mockStore.AssertNotCalled(t, "CompletePlanStep", mock.Anything, mock.Anything, mock.Anything) +} + func TestManager_GetAWSAccountID_Success(t *testing.T) { ctx := context.Background() mockSTS := new(MockSTSClient) diff --git a/internal/purchase/manager.go b/internal/purchase/manager.go index 971b68b7e..1ec3dd308 100644 --- a/internal/purchase/manager.go +++ b/internal/purchase/manager.go @@ -245,7 +245,7 @@ func (m *Manager) executeAndFinalize(ctx context.Context, exec *config.PurchaseE } } if execErr == nil { - if err := m.updatePlanProgress(ctx, exec.PlanID); err != nil { + if err := m.updatePlanProgress(ctx, exec); err != nil { logging.Errorf("Failed to update plan progress: %v", err) } } diff --git a/internal/purchase/manager_test.go b/internal/purchase/manager_test.go index 482a4f6fa..600b507b6 100644 --- a/internal/purchase/manager_test.go +++ b/internal/purchase/manager_test.go @@ -181,6 +181,7 @@ func TestManager_ProcessScheduledPurchases_DuePurchase(t *testing.T) { ExecutionID: "exec-123", PlanID: "plan-456", Status: "pending", + StepNumber: 1, ScheduledDate: pastDate, Recommendations: []config.RecommendationRecord{ { @@ -219,7 +220,7 @@ func TestManager_ProcessScheduledPurchases_DuePurchase(t *testing.T) { mockStore.On("SavePurchaseHistory", ctx, mock.AnythingOfType("*config.PurchaseHistoryRecord")).Return(nil) mockEmail.On("SendPurchaseConfirmation", ctx, mock.AnythingOfType("email.NotificationData")).Return(nil) mockStore.On("SavePurchaseExecution", ctx, mock.AnythingOfType("*config.PurchaseExecution")).Return(nil) - mockStore.On("IncrementPlanCurrentStep", ctx, "plan-456").Return(nil) + mockStore.On("CompletePlanStep", ctx, "plan-456", 1).Return(nil) mockSTS.On("GetCallerIdentity", ctx, mock.AnythingOfType("*sts.GetCallerIdentityInput")).Return(&sts.GetCallerIdentityOutput{ Account: aws.String("123456789012"), }, nil) @@ -596,6 +597,7 @@ func TestManager_RecoverStrandedApprovals_AWSOnlyRedrives(t *testing.T) { ExecutionID: "exec-aws-stranded", PlanID: "plan-aws-456", Status: "approved", + StepNumber: 1, Recommendations: []config.RecommendationRecord{ {Provider: "aws", Service: "ec2", ResourceType: "m5.large", Region: "us-east-1", Count: 1, UpfrontCost: 200.0, Selected: true, Purchased: false}, }, @@ -624,7 +626,7 @@ func TestManager_RecoverStrandedApprovals_AWSOnlyRedrives(t *testing.T) { mockStore.On("SavePurchaseExecution", ctx, mock.AnythingOfType("*config.PurchaseExecution")). Run(func(args mock.Arguments) { saved = args.Get(1).(*config.PurchaseExecution) }). Return(nil) - mockStore.On("IncrementPlanCurrentStep", ctx, "plan-aws-456").Return(nil) + mockStore.On("CompletePlanStep", ctx, "plan-aws-456", 1).Return(nil) mockSTS.On("GetCallerIdentity", ctx, mock.AnythingOfType("*sts.GetCallerIdentityInput")).Return(&sts.GetCallerIdentityOutput{ Account: aws.String("123456789012"), }, nil) @@ -684,6 +686,7 @@ func TestManager_RecoverStrandedApprovals_AzureReservationRedrives(t *testing.T) ExecutionID: "exec-azure-res-stranded", PlanID: "plan-azure-res", Status: "approved", + StepNumber: 1, Recommendations: []config.RecommendationRecord{ {Provider: "azure", Service: "compute", ResourceType: "Standard_D4s_v3", Region: "eastus", Count: 1, UpfrontCost: 300.0, Selected: true, Purchased: false}, }, @@ -712,7 +715,7 @@ func TestManager_RecoverStrandedApprovals_AzureReservationRedrives(t *testing.T) mockStore.On("SavePurchaseExecution", ctx, mock.AnythingOfType("*config.PurchaseExecution")). Run(func(args mock.Arguments) { saved = args.Get(1).(*config.PurchaseExecution) }). Return(nil) - mockStore.On("IncrementPlanCurrentStep", ctx, "plan-azure-res").Return(nil) + mockStore.On("CompletePlanStep", ctx, "plan-azure-res", 1).Return(nil) mockFactory.On("CreateAndValidateProvider", mock.Anything, "azure", mock.Anything).Return(mockProvider, nil) mockProvider.On("GetServiceClient", mock.Anything, common.ServiceCompute, "eastus").Return(mockServiceClient, nil) @@ -765,6 +768,7 @@ func TestManager_RecoverStrandedApprovals_GCPRedrives(t *testing.T) { ExecutionID: "exec-gcp-stranded", PlanID: "plan-gcp", Status: "approved", + StepNumber: 1, Recommendations: []config.RecommendationRecord{ {Provider: "gcp", Service: "compute", ResourceType: "n2-standard-4", Region: "us-central1", Count: 2, UpfrontCost: 150.0, Selected: true, Purchased: false}, }, @@ -793,7 +797,7 @@ func TestManager_RecoverStrandedApprovals_GCPRedrives(t *testing.T) { mockStore.On("SavePurchaseExecution", ctx, mock.AnythingOfType("*config.PurchaseExecution")). Run(func(args mock.Arguments) { saved = args.Get(1).(*config.PurchaseExecution) }). Return(nil) - mockStore.On("IncrementPlanCurrentStep", ctx, "plan-gcp").Return(nil) + mockStore.On("CompletePlanStep", ctx, "plan-gcp", 1).Return(nil) mockFactory.On("CreateAndValidateProvider", mock.Anything, "gcp", mock.Anything).Return(mockProvider, nil) mockProvider.On("GetServiceClient", mock.Anything, common.ServiceCompute, "us-central1").Return(mockServiceClient, nil) diff --git a/internal/purchase/money_path_regression_test.go b/internal/purchase/money_path_regression_test.go index 19f3714af..efe02beac 100644 --- a/internal/purchase/money_path_regression_test.go +++ b/internal/purchase/money_path_regression_test.go @@ -320,7 +320,7 @@ func TestMultiAccountPartialSuccessIsAcked(t *testing.T) { // updatePlanProgress only runs on a fully-clean run (execErr == nil); a // partial multi-account run skips it. Marked Maybe so a (non-deterministic) // all-success ordering of the two account goroutines doesn't fail the mock. - mockStore.On("IncrementPlanCurrentStep", ctx, "plan-x").Return(nil).Maybe() + mockStore.On("CompletePlanStep", ctx, "plan-x", mock.Anything).Return(nil).Maybe() var savedRoot *config.PurchaseExecution mockStore.SavePurchaseExecutionFn = func(_ context.Context, e *config.PurchaseExecution) error { @@ -412,7 +412,7 @@ func runScopelessRow(t *testing.T, exec *config.PurchaseExecution, accounts []co var tokens []string mockStore.On("SavePurchaseHistory", mock.Anything, mock.AnythingOfType("*config.PurchaseHistoryRecord")).Return(nil).Maybe() mockEmail.On("SendPurchaseConfirmation", mock.Anything, mock.AnythingOfType("email.NotificationData")).Return(nil).Maybe() - mockStore.On("IncrementPlanCurrentStep", mock.Anything, "plan-x").Return(nil).Maybe() + mockStore.On("CompletePlanStep", mock.Anything, "plan-x", mock.Anything).Return(nil).Maybe() mockFactory.On("CreateAndValidateProvider", mock.Anything, "aws", mock.Anything).Return(mockProviderInst, nil).Maybe() mockProviderInst.On("GetServiceClient", mock.Anything, common.ServiceEC2, mock.Anything).Return(mockServiceClient, nil).Maybe() mockServiceClient.On("PurchaseCommitment", mock.Anything, mock.Anything, mock.AnythingOfType("common.PurchaseOptions")). diff --git a/internal/purchase/notifications.go b/internal/purchase/notifications.go index f084a23b1..f9dede3df 100644 --- a/internal/purchase/notifications.go +++ b/internal/purchase/notifications.go @@ -128,10 +128,15 @@ func (m *Manager) getOrCreateExecution(ctx context.Context, plan *config.Purchas } tokenExpiresAt := time.Now().Add(config.ApprovalTokenTTL) execution := &config.PurchaseExecution{ - PlanID: plan.ID, - ExecutionID: uuid.New().String(), - Status: "pending", - StepNumber: plan.RampSchedule.CurrentStep, + PlanID: plan.ID, + ExecutionID: uuid.New().String(), + Status: "pending", + // step_number names the step this row will COMPLETE, not the count + // already completed, matching api.createPurchaseExecutionsTx + // (CurrentStep + i + 1). The ramp advance is keyed on this value since + // issue #1669, so stamping the completed count here would have every + // notification-created row re-complete a counted step, freezing the ramp. + StepNumber: plan.RampSchedule.CurrentStep + 1, ScheduledDate: *plan.NextExecutionDate, ApprovalToken: approvalToken, ApprovalTokenExpiresAt: &tokenExpiresAt, diff --git a/internal/purchase/notifications_test.go b/internal/purchase/notifications_test.go index ea8d352ca..d340913f7 100644 --- a/internal/purchase/notifications_test.go +++ b/internal/purchase/notifications_test.go @@ -191,7 +191,9 @@ func TestManager_GetOrCreateExecution(t *testing.T) { assert.NotNil(t, execution) assert.Equal(t, "plan-123", execution.PlanID) assert.Equal(t, "pending", execution.Status) - assert.Equal(t, 1, execution.StepNumber) + // step_number names the step this row will COMPLETE, so a plan with one + // step already done creates the row for step 2 (issue #1669). + assert.Equal(t, 2, execution.StepNumber) assert.NotEmpty(t, execution.ExecutionID) assert.NotEmpty(t, execution.ApprovalToken) @@ -325,7 +327,8 @@ func TestManager_GetOrCreateExecution_CreatesOnErrNotFound(t *testing.T) { require.NotNil(t, execution) assert.Equal(t, "plan-f2", execution.PlanID) assert.Equal(t, "pending", execution.Status) - assert.Equal(t, 2, execution.StepNumber) + // Two steps already completed, so this row is step 3 (issue #1669). + assert.Equal(t, 3, execution.StepNumber) assert.NotEmpty(t, execution.ExecutionID) assert.NotEmpty(t, execution.ApprovalToken) diff --git a/internal/purchase/ramp_step_progress_integration_test.go b/internal/purchase/ramp_step_progress_integration_test.go new file mode 100644 index 000000000..64f2e5e09 --- /dev/null +++ b/internal/purchase/ramp_step_progress_integration_test.go @@ -0,0 +1,443 @@ +//go:build integration +// +build integration + +package purchase + +// Regression coverage for issue #1669: retrying two failed accounts of one +// multi-account ramp step must advance the plan's CurrentStep exactly once. +// +// These tests run the real Manager against a real migrated Postgres (the plan +// row, the executions and the ramp advance all go through the production +// store), because the defect lives in the INTERACTION between the fan-out, the +// per-account retry successors and the store's step accounting. A unit test +// over the store helper alone cannot see it. + +import ( + "context" + "fmt" + "path/filepath" + "runtime" + "sync" + "testing" + "time" + + "github.com/LeanerCloud/CUDly/internal/config" + "github.com/LeanerCloud/CUDly/internal/database/postgres/migrations" + "github.com/LeanerCloud/CUDly/internal/database/postgres/testhelpers" + "github.com/LeanerCloud/CUDly/pkg/common" + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" +) + +// rampMigrationsPath resolves the migrations directory relative to this file. +func rampMigrationsPath() string { + _, filename, _, _ := runtime.Caller(0) + return filepath.Join(filepath.Dir(filename), "..", "database", "postgres", "migrations") +} + +// rampStepStore stands up a fresh Postgres container with all migrations +// applied and returns a production store bound to it, plus the pool so a test +// can drive its own transaction alongside the store's. +func rampStepStore(ctx context.Context, t *testing.T) (*config.PostgresStore, *pgxpool.Pool) { + t.Helper() + + container, err := testhelpers.SetupPostgresContainer(ctx, t) + require.NoError(t, err) + t.Cleanup(func() { + if cleanupErr := container.Cleanup(ctx); cleanupErr != nil { + t.Logf("container cleanup: %v", cleanupErr) + } + }) + + require.NoError(t, migrations.RunMigrations(ctx, container.DB.Pool(), rampMigrationsPath(), "", "")) + return config.NewPostgresStore(container.DB), container.DB.Pool() +} + +// rampStepFixture is the world a ramp-step test runs against: a plan sitting at +// CurrentStep=2 of 4 with three cloud accounts attached, and a manager wired to +// the real store with a provider that always sells. +type rampStepFixture struct { + store *config.PostgresStore + manager *Manager + plan *config.PurchasePlan + accounts []config.CloudAccount + + mu sync.Mutex + purchased []string // one entry per commitment that reached the fake cloud +} + +// newRampStepFixture builds the fixture. accountA is credential-resolvable from +// the start; accountB and accountC are configured with an auth mode that has no +// STS client wired, so their credential resolution fails deterministically -- +// modeling "two of three accounts failed this ramp step" without racing the +// provider mock. repairAccounts() flips them to a working auth mode, which is +// what an operator does before retrying. +// +// Enabled is false purely to keep the plan minimal: PurchasePlan.Validate +// requires a non-empty Services map only for enabled plans, and nothing on the +// execution path reads plan.Enabled. +func newRampStepFixture(ctx context.Context, t *testing.T) *rampStepFixture { + t.Helper() + + store, _ := rampStepStore(ctx, t) + + plan := &config.PurchasePlan{ + Name: "Ramp Step Plan", + AutoPurchase: true, + Services: map[string]config.ServiceConfig{}, + RampSchedule: config.RampSchedule{ + Type: "weekly", + PercentPerStep: 25, + StepIntervalDays: 7, + CurrentStep: 2, + TotalSteps: 4, + StartDate: time.Now().AddDate(0, 0, -14), + }, + } + require.NoError(t, store.CreatePurchasePlan(ctx, plan)) + + authModes := []string{"access_keys", "role_arn", "role_arn"} + accountIDs := make([]string, 0, len(authModes)) + for i, mode := range authModes { + acct := &config.CloudAccount{ + Name: fmt.Sprintf("acct-%c", 'a'+i), + Provider: "aws", + ExternalID: fmt.Sprintf("11111111111%d", i), + Enabled: true, + AWSAuthMode: mode, + AWSRoleARN: "arn:aws:iam::111111111110:role/cudly", + } + require.NoError(t, store.CreateCloudAccount(ctx, acct)) + accountIDs = append(accountIDs, acct.ID) + } + require.NoError(t, store.SetPlanAccounts(ctx, plan.ID, accountIDs)) + + accounts, err := store.GetPlanAccounts(ctx, plan.ID) + require.NoError(t, err) + require.Len(t, accounts, 3, "the fan-out must see all three accounts") + // GetPlanAccounts orders by name, so index 0 is the account that resolves + // and 1-2 are the ones that fail. Pin it: the whole scenario depends on + // which accounts fail, and a change in ordering would quietly retarget the + // retries at an account that never failed. + require.Equal(t, "access_keys", accounts[0].AWSAuthMode, "accounts[0] must be the account that succeeds") + require.Equal(t, "role_arn", accounts[1].AWSAuthMode, "accounts[1] must be a failing account") + require.Equal(t, "role_arn", accounts[2].AWSAuthMode, "accounts[2] must be a failing account") + + f := &rampStepFixture{store: store, plan: plan, accounts: accounts} + + mockEmail := new(MockEmailSender) + mockFactory := new(MockProviderFactory) + mockProviderInst := new(MockProvider) + mockServiceClient := new(MockServiceClient) + t.Cleanup(func() { mockEmail.AssertExpectations(t) }) + t.Cleanup(func() { mockFactory.AssertExpectations(t) }) + t.Cleanup(func() { mockProviderInst.AssertExpectations(t) }) + t.Cleanup(func() { mockServiceClient.AssertExpectations(t) }) + + mockEmail.On("SendPurchaseConfirmation", mock.Anything, mock.AnythingOfType("email.NotificationData")).Return(nil).Maybe() + mockFactory.On("CreateAndValidateProvider", mock.Anything, "aws", mock.Anything).Return(mockProviderInst, nil).Maybe() + mockProviderInst.On("GetServiceClient", mock.Anything, common.ServiceEC2, mock.Anything).Return(mockServiceClient, nil).Maybe() + mockServiceClient.On("PurchaseCommitment", mock.Anything, mock.Anything, mock.AnythingOfType("common.PurchaseOptions")). + Run(func(args mock.Arguments) { + opts := args.Get(2).(common.PurchaseOptions) + f.mu.Lock() + f.purchased = append(f.purchased, opts.IdempotencyToken) + f.mu.Unlock() + }). + Return(common.PurchaseResult{Success: true, CommitmentID: "ri-ok"}, nil).Maybe() + + f.manager = &Manager{ + config: store, + email: mockEmail, + providerFactory: mockFactory, + credStore: awsAccessKeyCredStore(), + dashboardURL: "https://dashboard.example.com", + } + return f +} + +// repairAccounts flips the two broken accounts to a resolvable auth mode, the +// operator-side precondition of a successful retry. +func (f *rampStepFixture) repairAccounts(ctx context.Context, t *testing.T) { + t.Helper() + for i := 1; i < len(f.accounts); i++ { + acct := f.accounts[i] + acct.AWSAuthMode = "access_keys" + require.NoError(t, f.store.UpdateCloudAccount(ctx, &acct)) + } +} + +// currentStep re-reads the plan from Postgres and returns its ramp position. +func (f *rampStepFixture) currentStep(ctx context.Context, t *testing.T) int { + t.Helper() + plan, err := f.store.GetPurchasePlan(ctx, f.plan.ID) + require.NoError(t, err) + return plan.RampSchedule.CurrentStep +} + +// purchaseCount reports how many commitments reached the fake cloud so far. +func (f *rampStepFixture) purchaseCount() int { + f.mu.Lock() + defer f.mu.Unlock() + return len(f.purchased) +} + +// rampStepRecommendation is the single rec every execution in these tests buys. +func rampStepRecommendation() []config.RecommendationRecord { + return []config.RecommendationRecord{ + {Provider: "aws", Service: "ec2", ResourceType: "m5.large", Region: "us-east-1", Count: 1, UpfrontCost: 300, Selected: true}, + } +} + +// saveRootExecution persists the plan-scheduled root row for ramp step +// stepNumber: no cloud_account_id, so the executor fans it out. +func (f *rampStepFixture) saveRootExecution(ctx context.Context, t *testing.T, stepNumber int) *config.PurchaseExecution { + t.Helper() + exec := &config.PurchaseExecution{ + ExecutionID: uuid.New().String(), + IdempotencyKey: fmt.Sprintf("lineage-step-%d", stepNumber), + PlanID: f.plan.ID, + Status: "pending", + StepNumber: stepNumber, + ScheduledDate: time.Now(), + Recommendations: rampStepRecommendation(), + } + require.NoError(t, f.store.SavePurchaseExecution(ctx, exec)) + return exec +} + +// saveRetryExecution persists the successor row an operator's Retry produces +// for one failed per-account row, mirroring api.persistRetryExecution: the +// PlanID, StepNumber, CloudAccountID and idempotency lineage all propagate from +// the predecessor, and the row arrives already approved by the human who +// clicked Retry. +func (f *rampStepFixture) saveRetryExecution(ctx context.Context, t *testing.T, stepNumber int, accountID string) *config.PurchaseExecution { + t.Helper() + acctID := accountID + exec := &config.PurchaseExecution{ + ExecutionID: uuid.New().String(), + IdempotencyKey: fmt.Sprintf("lineage-step-%d:%s", stepNumber, accountID), + PlanID: f.plan.ID, + CloudAccountID: &acctID, + Status: "approved", + StepNumber: stepNumber, + ScheduledDate: time.Now(), + Recommendations: rampStepRecommendation(), + Source: common.PurchaseSourceWeb, + RetryAttemptN: 1, + } + require.NoError(t, f.store.SavePurchaseExecution(ctx, exec)) + return exec +} + +// execute drives the real async executor entry point for one row. +func (f *rampStepFixture) execute(ctx context.Context, exec *config.PurchaseExecution) error { + msg := AsyncMessage{Type: MessageTypeExecutePurchase, ExecutionID: exec.ExecutionID} + return f.manager.handleExecutePurchase(ctx, &msg) +} + +// TestRampStepAdvancesOncePerStepAcrossPerAccountRetries is the issue #1669 +// regression guard. +// +// A multi-account plan fans ramp step 3 out over accounts A, B and C. A buys; B +// and C fail on credentials. The operator repairs both and retries each failed +// row, and both retries succeed. The ramp must end on step 3, not step 4: with +// the pre-fix blind CurrentStep++ each successful retry advanced the ramp +// again, so one whole ramp step's worth of commitment silently dropped out of +// the plan's accounting. +func TestRampStepAdvancesOncePerStepAcrossPerAccountRetries(t *testing.T) { + ctx := context.Background() + f := newRampStepFixture(ctx, t) + + require.Equal(t, 2, f.currentStep(ctx, t), "fixture must start on step 2 of 4") + + // Step 3 fans out: A commits, B and C fail credential resolution. + root := f.saveRootExecution(ctx, t, 3) + require.NoError(t, f.execute(ctx, root), + "a partial multi-account run must be acked, not surfaced as a flat failure (#1014)") + + require.Equal(t, 1, f.purchaseCount(), "exactly one account (A) may commit on the first pass") + assert.Equal(t, 2, f.currentStep(ctx, t), + "a partially-failed step must not advance the ramp at all") + + // The operator fixes the two broken accounts and retries each failed row. + f.repairAccounts(ctx, t) + + retryB := f.saveRetryExecution(ctx, t, 3, f.accounts[1].ID) + require.NoError(t, f.execute(ctx, retryB)) + require.Equal(t, 2, f.purchaseCount(), "B's retry must commit") + assert.Equal(t, 3, f.currentStep(ctx, t), + "the first successful completion of step 3 advances the ramp to 3") + + retryC := f.saveRetryExecution(ctx, t, 3, f.accounts[2].ID) + require.NoError(t, f.execute(ctx, retryC)) + require.Equal(t, 3, f.purchaseCount(), "C's retry must commit") + + assert.Equal(t, 3, f.currentStep(ctx, t), + "completing step 3 a second time must NOT advance the ramp again (issue #1669): "+ + "the plan has bought 3 of 4 ramp steps and must still say so") + + // Completing step 3 once more against this exact post-scenario state must + // take the no-op branch. Nothing else observable distinguishes it from the + // refusal branch (both return before the write and both leave CurrentStep + // at 3), and the caller only logs the error, so assert on the return value + // directly against real DB state rather than against pgxmock. + require.NoError(t, f.store.CompletePlanStep(ctx, f.plan.ID, 3), + "a counted step must no-op, not be refused as a skipped predecessor") + + // The ramp is not complete, so the plan still points at a next execution. + plan, err := f.store.GetPurchasePlan(ctx, f.plan.ID) + require.NoError(t, err) + assert.False(t, plan.RampSchedule.IsComplete(), "3 of 4 steps bought is not a complete ramp") + require.NotNil(t, plan.NextExecutionDate, "an incomplete ramp must keep a next execution date") + + // And the 4th step still counts. Pre-fix the ramp was already sitting at 4 + // here, so this step's completion changed nothing and the plan finished + // having bought three steps while reporting four. + step4 := f.saveRootExecution(ctx, t, 4) + require.NoError(t, f.execute(ctx, step4)) + assert.Equal(t, 4, f.currentStep(ctx, t), "step 4 completes the ramp") + assert.Equal(t, 6, f.purchaseCount(), "step 4 commits once per account") +} + +// TestCompletePlanStepIsIdempotentUnderConcurrency answers the question the +// row lock is supposed to settle: can two concurrent completions of the SAME +// ramp step both pass the guard and both advance? +// +// The goroutines are released together from a shared start gate and contend on +// the plan row for real (SELECT ... FOR UPDATE inside the store's transaction), +// so this is a measurement, not an argument from isolation levels. +func TestCompletePlanStepIsIdempotentUnderConcurrency(t *testing.T) { + ctx := context.Background() + store, _ := rampStepStore(ctx, t) + plan := saveConcurrencyRampPlan(ctx, t, store, "Concurrent Ramp Plan") + + const racers = 8 + start := make(chan struct{}) + errCh := make(chan error, racers) + var wg sync.WaitGroup + for i := 0; i < racers; i++ { + wg.Add(1) + go func() { + defer wg.Done() + <-start + errCh <- store.CompletePlanStep(ctx, plan.ID, 3) + }() + } + close(start) + wg.Wait() + close(errCh) + + seen := 0 + for err := range errCh { + seen++ + require.NoError(t, err, "every concurrent completion of step 3 must succeed or no-op, never error") + } + require.Equal(t, racers, seen, "all racers must have reported") + + reloaded, err := store.GetPurchasePlan(ctx, plan.ID) + require.NoError(t, err) + assert.Equal(t, 3, reloaded.RampSchedule.CurrentStep, + "%d concurrent completions of step 3 must leave the ramp on step 3", racers) +} + +// saveConcurrencyRampPlan persists a plan sitting on step 2 of 4. +func saveConcurrencyRampPlan(ctx context.Context, t *testing.T, store *config.PostgresStore, name string) *config.PurchasePlan { + t.Helper() + plan := &config.PurchasePlan{ + Name: name, + Services: map[string]config.ServiceConfig{}, + RampSchedule: config.RampSchedule{ + Type: "weekly", PercentPerStep: 25, StepIntervalDays: 7, + CurrentStep: 2, TotalSteps: 4, StartDate: time.Now().AddDate(0, 0, -14), + }, + } + require.NoError(t, store.CreatePurchasePlan(ctx, plan)) + return plan +} + +// TestCompletePlanStepBlocksOnTheRowLockAndSeesTheWinner forces the interleaving +// the previous test can only hope for, so the idempotency guard's dependence on +// the row lock is measured rather than argued. +// +// A transaction outside the store takes the plan row's FOR UPDATE lock and +// advances the ramp to step 3, exactly as the first per-account retry's +// transaction would. While that transaction is open, a concurrent +// CompletePlanStep(plan, 3) -- the second retry -- must block rather than read +// the pre-advance value: if it could proceed, both would see CurrentStep=2 and +// both would advance, which is the lost update the lock exists to prevent. +// Once the holder commits, the blocked call must observe the committed step 3 +// and no-op. +func TestCompletePlanStepBlocksOnTheRowLockAndSeesTheWinner(t *testing.T) { + ctx := context.Background() + store, pool := rampStepStore(ctx, t) + plan := saveConcurrencyRampPlan(ctx, t, store, "Locked Ramp Plan") + + holder, err := pool.Begin(ctx) + require.NoError(t, err) + defer func() { _ = holder.Rollback(ctx) }() + + var lockedID string + require.NoError(t, + holder.QueryRow(ctx, `SELECT id FROM purchase_plans WHERE id = $1 FOR UPDATE`, plan.ID).Scan(&lockedID)) + require.Equal(t, plan.ID, lockedID) + + // The winning retry's advance, still uncommitted. + _, err = holder.Exec(ctx, + `UPDATE purchase_plans SET ramp_schedule = jsonb_set(ramp_schedule, '{current_step}', '3') WHERE id = $1`, + plan.ID) + require.NoError(t, err) + + done := make(chan error, 1) + go func() { done <- store.CompletePlanStep(ctx, plan.ID, 3) }() + + // Wait until Postgres reports the completion's own SELECT ... FOR UPDATE + // blocked on a lock, so the "still running" check below is not a race + // against goroutine start-up. Matching the query text matters: if the store + // stopped taking the row lock on the READ, the blocked statement would be + // the later UPDATE instead and this poll would never match. + require.Eventually(t, func() bool { + var blocked int + if qErr := pool.QueryRow(ctx, ` + SELECT count(*) FROM pg_stat_activity + WHERE wait_event_type = 'Lock' + AND state = 'active' + AND pid <> pg_backend_pid() + AND query ILIKE '%purchase_plans%FOR UPDATE%'`, + ).Scan(&blocked); qErr != nil { + return false + } + return blocked > 0 + }, 15*time.Second, 25*time.Millisecond, + "the concurrent completion's locking read must wait on the plan row lock") + + select { + case earlyErr := <-done: + t.Fatalf("CompletePlanStep returned (%v) while another transaction held the plan row lock: it did not serialize", earlyErr) + default: + } + + require.NoError(t, holder.Commit(ctx)) + + select { + case loserErr := <-done: + require.NoError(t, loserErr, "the loser of the race must no-op cleanly, not error") + case <-time.After(15 * time.Second): + t.Fatal("CompletePlanStep never returned after the row lock was released") + } + + reloaded, err := store.GetPurchasePlan(ctx, plan.ID) + require.NoError(t, err) + assert.Equal(t, 3, reloaded.RampSchedule.CurrentStep, + "the blocked completion must observe the committed step 3 and leave it alone") + // CurrentStep alone cannot tell the two outcomes apart: a stale read of + // step 2 followed by an advance also lands on 3. last_execution_date can. + // The holder's UPDATE never touched it, and only the advancing path in + // CompletePlanStep writes it, so its absence proves the blocked call took + // the no-op branch rather than re-advancing off a pre-commit read. + assert.Nil(t, reloaded.LastExecutionDate, + "the blocked completion must no-op, not re-advance the ramp from a stale read") +} diff --git a/internal/server/test_helpers_test.go b/internal/server/test_helpers_test.go index f1ee2611f..ce0afc499 100644 --- a/internal/server/test_helpers_test.go +++ b/internal/server/test_helpers_test.go @@ -49,7 +49,7 @@ func (m *mockConfigStoreForHealth) UpdatePurchasePlan(ctx context.Context, plan return nil } -func (m *mockConfigStoreForHealth) IncrementPlanCurrentStep(_ context.Context, _ string) error { +func (m *mockConfigStoreForHealth) CompletePlanStep(_ context.Context, _ string, _ int) error { return nil } From 6c8c25d42e7b51525eac326b39662806821babd7 Mon Sep 17 00:00:00 2001 From: Cristian Magherusan-Stanciu Date: Wed, 19 Aug 2026 15:52:14 +0200 Subject: [PATCH 2/2] fix(plans): record a refused ramp advance on the execution row The refusal paths are new in #1669: the previous blind CurrentStep++ could not decline, so a purchase always moved the ramp. CompletePlanStep can now decline (an unknown step_number, or a step more than one beyond CurrentStep), which means money can be spent on a step the plan then declines to count. That stall is not self-correcting. CompletePlanStep returns before its write, so next_execution_date stays stale, and shouldNotifyPlan reads a stale date as daysUntil < 0 and stops notifying the plan entirely. A plan on the notification-driven path therefore goes quiet rather than retrying, and the only trace was a logging.Errorf. Stamp the refusal on the execution row that just completed so it outlives the log retention window and surfaces in History. The status stays completed: the purchase did complete, and only the progress accounting did not. Recovery is deliberately not scheduled here; deriving ramp progress from the executions table is tracked in #1861. --- internal/purchase/execution_test.go | 77 +++++++++++++++++++++++++++++ internal/purchase/manager.go | 24 +++++++++ 2 files changed, 101 insertions(+) diff --git a/internal/purchase/execution_test.go b/internal/purchase/execution_test.go index 67e79549a..e4f71b544 100644 --- a/internal/purchase/execution_test.go +++ b/internal/purchase/execution_test.go @@ -444,6 +444,83 @@ func TestManager_UpdatePlanProgress_ZeroStepRefusesAndSkipsStore(t *testing.T) { mockStore.AssertNotCalled(t, "CompletePlanStep", mock.Anything, mock.Anything, mock.Anything) } +// TestExecuteAndFinalize_RefusedRampAdvanceIsRecordedOnTheRow is the CR +// follow-up guard on #1862: the refusal paths are new in issue #1669 (the old +// blind CurrentStep++ could not decline), so money can now be spent on a step +// the plan then declines to count. Before this, the only trace was a log line, +// and a stalled plan is not self-correcting because shouldNotifyPlan reads the +// stale next_execution_date as daysUntil < 0 and stops notifying. +// +// The execution must stay "completed" (the purchase did complete) while the +// refusal is persisted on the row so History shows it. +func TestExecuteAndFinalize_RefusedRampAdvanceIsRecordedOnTheRow(t *testing.T) { + ctx := context.Background() + mockStore := new(MockConfigStore) + mockEmail := new(MockEmailSender) + mockFactory := new(MockProviderFactory) + mockProviderInst := new(MockProvider) + mockServiceClient := new(MockServiceClient) + t.Cleanup(func() { mockStore.AssertExpectations(t) }) + + // A plan-attributed row carrying no ramp step: the deploy-overlap shape that + // updatePlanProgress refuses. CompletePlanStep is registered but must never + // be reached, so the refusal comes from the caller rather than the store. + exec := &config.PurchaseExecution{ + ExecutionID: "exec-refused-advance", + PlanID: "plan-refuse", + Status: "running", + StepNumber: 0, + Recommendations: []config.RecommendationRecord{ + {Provider: "aws", Service: "ec2", ResourceType: "m5.large", Region: "us-east-1", Count: 1, UpfrontCost: 300, Selected: true}, + }, + } + + var saves []config.PurchaseExecution + mockStore.SavePurchaseExecutionFn = func(_ context.Context, e *config.PurchaseExecution) error { + saves = append(saves, *e) + return nil + } + mockStore.GetPurchasePlanFn = func(_ context.Context, id string) (*config.PurchasePlan, error) { + return &config.PurchasePlan{ID: id, Name: "Refuse Plan"}, nil + } + mockStore.GetPlanAccountsFn = func(_ context.Context, _ string) ([]config.CloudAccount, error) { + return nil, nil + } + mockStore.On("CompletePlanStep", mock.Anything, mock.Anything, mock.Anything).Return(nil).Maybe() + mockStore.On("SavePurchaseHistory", mock.Anything, mock.AnythingOfType("*config.PurchaseHistoryRecord")).Return(nil).Maybe() + mockEmail.On("SendPurchaseConfirmation", mock.Anything, mock.AnythingOfType("email.NotificationData")).Return(nil).Maybe() + mockFactory.On("CreateAndValidateProvider", mock.Anything, "aws", mock.Anything).Return(mockProviderInst, nil).Maybe() + mockProviderInst.On("GetServiceClient", mock.Anything, common.ServiceEC2, mock.Anything).Return(mockServiceClient, nil).Maybe() + mockServiceClient.On("PurchaseCommitment", mock.Anything, mock.Anything, mock.Anything). + Return(common.PurchaseResult{Success: true, CommitmentID: "ri-ok"}, nil).Maybe() + + manager := &Manager{ + config: mockStore, + email: mockEmail, + providerFactory: mockFactory, + credStore: awsAccessKeyCredStore(), + dashboardURL: "https://dashboard.example.com", + } + + require.NoError(t, manager.executeAndFinalize(ctx, exec), + "a refused ramp advance must not fail the purchase; the money already moved") + + // Assert a non-zero number of saves before asserting anything about their + // content, so an empty slice cannot satisfy the checks below. + require.GreaterOrEqual(t, len(saves), 2, + "the terminal save plus the ramp-refusal save must both reach the store") + + final := saves[len(saves)-1] + assert.Equal(t, "completed", final.Status, + "the purchase completed; only the plan's progress accounting did not") + assert.Contains(t, final.Error, "ramp not advanced", + "the refusal must survive on the row, not only in the log") + assert.Contains(t, final.Error, "plan-refuse", + "the recorded note must name the plan whose ramp stalled") + + mockStore.AssertNotCalled(t, "CompletePlanStep", mock.Anything, mock.Anything, mock.Anything) +} + func TestManager_GetAWSAccountID_Success(t *testing.T) { ctx := context.Background() mockSTS := new(MockSTSClient) diff --git a/internal/purchase/manager.go b/internal/purchase/manager.go index 1ec3dd308..a30f76b48 100644 --- a/internal/purchase/manager.go +++ b/internal/purchase/manager.go @@ -247,11 +247,35 @@ func (m *Manager) executeAndFinalize(ctx context.Context, exec *config.PurchaseE if execErr == nil { if err := m.updatePlanProgress(ctx, exec); err != nil { logging.Errorf("Failed to update plan progress: %v", err) + m.recordRampAdvanceRefusal(ctx, exec, err) } } return execErr } +// recordRampAdvanceRefusal stamps a refused ramp advance onto the execution row +// that just completed, so the decision outlives the log retention window. +// +// The refusal paths are new in issue #1669: the previous blind CurrentStep++ +// could not decline, so a purchase always moved the ramp. Now money can be spent +// on a step the plan then declines to count (an unknown step_number, or a step +// more than one beyond CurrentStep), and CompletePlanStep returning before its +// write leaves next_execution_date stale, which shouldNotifyPlan reads as +// daysUntil < 0 and stops notifying that plan. A stall is therefore not +// self-correcting, and a logging.Errorf was its only trace. +// +// The status stays "completed": the purchase did complete, and only the plan's +// progress accounting did not. Recovery is deliberately NOT scheduled here (see +// issue #1861); this records the fact so History shows it and an operator can +// act on it. +func (m *Manager) recordRampAdvanceRefusal(ctx context.Context, exec *config.PurchaseExecution, cause error) { + exec.Error = appendErrNote(exec.Error, fmt.Sprintf("ramp not advanced: %v", cause)) + if saveErr := m.config.SavePurchaseExecution(ctx, exec); saveErr != nil { + logging.Errorf("AUDIT LOSS: failed to persist ramp-advance refusal for execution %s: %v", + exec.ExecutionID, saveErr) + } +} + // allRecsSafeToRedrive reports whether this sweep's automatic in-place re-drive // may run for exec. Two conditions: every recommendation must be safe to // re-drive, which RedriveRefusalReason below owns and documents, AND the