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..e4f71b544 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,125 @@ 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) +} + +// 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 971b68b7e..a30f76b48 100644 --- a/internal/purchase/manager.go +++ b/internal/purchase/manager.go @@ -245,13 +245,37 @@ 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) + 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 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 }