Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion internal/analytics/collector_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down
2 changes: 1 addition & 1 deletion internal/api/handler_purchases_retry_fanout_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
4 changes: 2 additions & 2 deletions internal/api/handler_purchases_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
14 changes: 10 additions & 4 deletions internal/config/interfaces.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
54 changes: 45 additions & 9 deletions internal/config/store_postgres.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -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() {
Expand All @@ -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
Expand Down
Loading
Loading