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
11 changes: 10 additions & 1 deletion internal/config/errors.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,8 +7,17 @@ var ErrNotFound = errors.New("not found")

// ErrExecutionNotInExpectedStatus is returned by TransitionExecutionStatus
// when the target execution exists but its current status is not in the
// allowed `fromStatuses` set — i.e. the atomic CAS rejected because some
// allowed `fromStatuses` set -- i.e. the atomic CAS rejected because some
// other writer transitioned the row first (e.g. the real executor finished
// between the reaper's SELECT and CAS). Callers can use errors.Is to
// distinguish this legitimate race-loss from a hard DB error.
var ErrExecutionNotInExpectedStatus = errors.New("execution not in expected status")

// ErrAuditLoss is returned (wrapped) by executeAndFinalize when the purchase
// run itself completed but the subsequent SavePurchaseExecution call failed.
// The execution is already "running" (per the CAS in claimAndRedrive) but its
// final state was never persisted -- the row is stranded in "running" until
// the next recovery sweep. Callers that silence all drive errors (e.g.
// claimAndRedrive) must propagate this sentinel so the sweep surfaces the
// persistence failure rather than silently dropping the stranded row.
var ErrAuditLoss = errors.New("audit loss: execution persistence failed after purchase")
226 changes: 185 additions & 41 deletions internal/purchase/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -162,8 +162,18 @@ func (m *Manager) executeAndFinalize(ctx context.Context, exec *config.PurchaseE
if !wasMultiAccount {
if err := m.config.SavePurchaseExecution(ctx, exec); err != nil {
logging.Errorf("AUDIT LOSS: failed to save execution status: %v", err)
if execErr == nil {
execErr = fmt.Errorf("audit loss: %w", err)
// Wrap with ErrAuditLoss regardless of whether executePurchase itself
// failed. When execErr != nil (provider/partial error), finalizeExecution
// stamped a terminal status on the in-memory exec struct, but if
// SavePurchaseExecution then failed the DB row is still in "running" --
// exactly the stranded-row scenario ErrAuditLoss signals. Preserve the
// original execErr as the innermost %w so errors.As/errors.Is can still
// reach it from callers (e.g. claimAndRedrive checking ErrAuditLoss).
if execErr != nil {
execErr = fmt.Errorf("%w: terminal save failed (%v); original execution error: %w",
config.ErrAuditLoss, err, execErr)
} else {
execErr = fmt.Errorf("%w: %w", config.ErrAuditLoss, err)
}
}
}
Expand All @@ -175,22 +185,161 @@ func (m *Manager) executeAndFinalize(ctx context.Context, exec *config.PurchaseE
return execErr
}

// allRecsAWS reports whether every recommendation in the execution targets AWS.
// Empty provider ("") is treated as AWS because pre-multi-cloud rows predate the
// provider field and are all AWS. An execution with no recommendations returns
// false so it falls through to the safe-fail path (nothing to re-drive anyway).
//
// Azure and GCP recs are excluded: Azure re-drive idempotency is blocked by issue
// #721 (#639), and GCP commitments are currently flagged unsupported (#640). Only
// call executeAndFinalize on an execution when this predicate returns true AND the
// execution has a non-empty ExecutionID (needed for DeriveIdempotencyToken).
func allRecsAWS(exec *config.PurchaseExecution) bool {
if len(exec.Recommendations) == 0 {
return false
}
for _, rec := range exec.Recommendations {
if rec.Provider != "" && rec.Provider != "aws" {
return false
}
}
return true
}

// claimAndRedrive atomically claims a stranded AWS-only execution (by
// transitioning its status from "approved" to "running") and then re-drives it
// via executeAndFinalize. It is extracted from RecoverStrandedApprovals to keep
// that function's cyclomatic complexity within the gocyclo:10 limit.
//
// Returns (true, nil) when the claim was won and the re-drive completed (row is
// now in a terminal state). A re-drive error is not fatal -- executeAndFinalize
// already stamped a terminal status -- so it returns (false, nil) on drive
// failure too (the failed row is still visible in History).
Comment thread
coderabbitai[bot] marked this conversation as resolved.
// Returns (false, nil) when the CAS claim is lost to a concurrent sweep or the
// row vanishes mid-flight (both are benign races).
// Returns (false, err) only on a genuine DB error during the claim step.
func (m *Manager) claimAndRedrive(ctx context.Context, exec *config.PurchaseExecution) (bool, error) {
// Atomically claim ownership before re-driving to prevent two concurrent
// sweeps (or a late original completion) from both calling executeAndFinalize
// on the same approved row. The CAS transitions "approved" -> "running";
// only the winner proceeds.
claimed, claimErr := m.config.TransitionExecutionStatus(ctx, exec.ExecutionID, []string{"approved"}, "running")
if claimErr != nil {
// ErrNotFound: row vanished between SELECT and CAS - benign race.
// ErrExecutionNotInExpectedStatus: another sweep or the original run
// already claimed/completed this row - also benign.
if errors.Is(claimErr, config.ErrNotFound) || errors.Is(claimErr, config.ErrExecutionNotInExpectedStatus) {
logging.Warnf("Skipping re-drive of %s (CAS claim lost to concurrent worker): %v", exec.ExecutionID, claimErr)
return false, nil
}
return false, fmt.Errorf("failed to claim execution %s for re-drive: %w", exec.ExecutionID, claimErr)
}
// Update the local struct to reflect the committed DB state so
// finalizeExecution starts from "running" rather than "approved".
exec.Status = claimed.Status
logging.Infof("Recovering stranded AWS-only execution %s via idempotent re-drive (issue #632)", exec.ExecutionID)
if driveErr := m.executeAndFinalize(ctx, exec); driveErr != nil {
logging.Errorf("Re-drive of stranded execution %s failed: %v", exec.ExecutionID, driveErr)
// Persistence failures (ErrAuditLoss) are non-benign: the row was CAS-ed
// to "running" but SavePurchaseExecution failed, so no terminal status was
// persisted. Propagate so the sweep surfaces the error rather than silently
// dropping a row that is now stranded in "running".
if errors.Is(driveErr, config.ErrAuditLoss) {
return false, fmt.Errorf("persistence failure re-driving execution %s (row stranded in running): %w", exec.ExecutionID, driveErr)
}
// Benign provider/rec errors: finalizeExecution already stamped a terminal
// status (failed/partially_completed) and SavePurchaseExecution succeeded.
// The row is in a terminal state; log and continue without counting as
// recovered.
return false, nil
}
return true, nil
}

// safeFail atomically transitions a stranded execution to "failed" and stamps a
// recovery error on it. It is extracted from RecoverStrandedApprovals to keep
// that function's cyclomatic complexity within the gocyclo:10 limit.
//
// Returns (true, nil) when the row was successfully transitioned to "failed".
// Returns (false, nil) when TransitionExecutionStatus fails but the row has
// already left "approved" (benign race - the original run completed late;
// not counted as a recovery since no action was taken here).
// Returns (false, err) when a real store failure occurs.
func (m *Manager) safeFail(ctx context.Context, exec *config.PurchaseExecution) (bool, error) {
logging.Errorf("Recovering stranded approved execution %s (approved but never finalized; failing it for visibility)", exec.ExecutionID)

updated, txErr := m.config.TransitionExecutionStatus(ctx, exec.ExecutionID, []string{"approved"}, "failed")
if txErr != nil {
// ErrNotFound means the row vanished between the stale SELECT and
// this CAS attempt (e.g. deleted by an operator or a concurrent
// sweep already claimed and deleted it). That is a benign race-loss:
// there is nothing left to fail, and the caller should not be
// charged with an error.
if errors.Is(txErr, config.ErrNotFound) {
logging.Warnf("Skipping recovery of %s (row no longer exists, benign race-loss): %v", exec.ExecutionID, txErr)
return false, nil
}
// ErrExecutionNotInExpectedStatus means the row exists but its
// status has already moved out of "approved" (a concurrent sweep or
// the original run won the CAS race). Treat identically to the
// ErrNotFound case: nothing left to do here, no re-read needed.
// This mirrors claimAndRedrive and reaper.go, which both treat this
// sentinel as terminally benign.
if errors.Is(txErr, config.ErrExecutionNotInExpectedStatus) {
logging.Warnf("Skipping recovery of %s (row already left approved state, benign CAS race-loss): %v", exec.ExecutionID, txErr)
return false, nil
}
// Distinguish benign races (row already left the "approved"
// state - concurrent sweep handled it, or the original run
// finished after the LIST snapshot) from real store
// failures (DB unreachable, query syntax error). A real
// store failure must fail the sweep so a transient DB
// outage does not silently under-recover. We probe the
// current row state via GetExecutionByID: a clean read
// with Status != "approved" confirms the race; any other
// outcome (read error, still-approved row) is a real
// failure worth propagating.
current, getErr := m.config.GetExecutionByID(ctx, exec.ExecutionID)
if getErr == nil && current != nil && current.Status != "approved" {
logging.Warnf("Skipping recovery of %s (already transitioned out of approved): %v", exec.ExecutionID, txErr)
return false, nil
}
return false, fmt.Errorf("failed to transition stranded execution %s to failed: %w", exec.ExecutionID, txErr)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

updated.Error = "execution was approved but its purchase run was interrupted before completing and never finalized; failed by the recovery sweep so it is not silently stuck (issue #632). Verify on the cloud provider that no commitment was created, then Retry."
if saveErr := m.config.SavePurchaseExecution(ctx, updated); saveErr != nil {
// The atomic flip to "failed" already landed via TransitionExecutionStatus;
// only the explanatory error string failed to persist. Log loudly but
// still count the recovery - the row is no longer stranded in "approved".
logging.Errorf("AUDIT GAP: failed to stamp recovery error on %s: %v", exec.ExecutionID, saveErr)
}
return true, nil
}

// RecoverStrandedApprovals finds executions stuck in the "approved" status past
// staleApprovedThreshold and drives them into a terminal "failed" state so an
// approved row can never sit permanently stranded with no owner (issue #632).
// staleApprovedThreshold and either re-drives them idempotently (AWS-only
// executions with a durable ExecutionID) or drives them into a terminal "failed"
// state (mixed/Azure/GCP or legacy rows without a stable ExecutionID).
//
// It deliberately does NOT re-run the purchase: there is no idempotency token on
// commitment creation (EC2 PurchaseReservedInstancesOffering sets no ClientToken;
// CreateSavingsPlan has none), so an automatic re-drive of a row that was
// interrupted *after* AWS created the commitment but *before* the row persisted
// would double-purchase. Failing the row makes it visible in the History view
// (which surfaces failed rows) and Retry-able by an operator who has confirmed
// the AWS-side state, instead of requiring a manual DB edit.
// AWS-only path (issue #632 Option 5): all AWS executors derive a deterministic
// per-rec idempotency token via common.DeriveIdempotencyToken(exec.ExecutionID, i)
// at purchase time (execution.go:428). Re-driving with the same ExecutionID
// produces the same token, so AWS dedupes the second call and no double-purchase
// occurs. The row transitions from "approved" directly to "completed" (or
// "failed"/"partially_completed" on a genuine error), bypassing the manual Retry
// step required by the old safe-fail path.
//
// The transition is atomic: TransitionExecutionStatus only flips rows still in
// "approved", so if the original run finally completes between the stale SELECT
// and this UPDATE, the transition is a no-op and the genuine "completed" status
// is preserved.
// Safe-fail path (mixed/Azure/GCP/legacy): Azure two-step idempotency is blocked
// by issue #721; GCP is flagged unsupported (#640). Executions without a stable
// ExecutionID cannot derive tokens safely. These fall through to the original
// behaviour: the row is atomically transitioned to "failed" so it surfaces in
// History and can be Retry-ed by an operator after confirming the cloud-side state.
//
// The transition in the safe-fail path is atomic: TransitionExecutionStatus only
// flips rows still in "approved", so if the original run finally completes between
// the stale SELECT and this UPDATE, the transition is a no-op and the genuine
// "completed" status is preserved.
func (m *Manager) RecoverStrandedApprovals(ctx context.Context) (int, error) {
stranded, err := m.config.GetStaleApprovedExecutions(ctx, staleApprovedThreshold)
if err != nil {
Expand All @@ -200,36 +349,31 @@ func (m *Manager) RecoverStrandedApprovals(ctx context.Context) (int, error) {
recovered := 0
for i := range stranded {
exec := &stranded[i]
logging.Errorf("Recovering stranded approved execution %s (approved but never finalized; failing it for visibility)", exec.ExecutionID)

updated, txErr := m.config.TransitionExecutionStatus(ctx, exec.ExecutionID, []string{"approved"}, "failed")
if txErr != nil {
// Distinguish benign races (row already left the "approved"
// state — concurrent sweep handled it, or the original run
// finished after the LIST snapshot) from real store
// failures (DB unreachable, query syntax error). A real
// store failure must fail the sweep so a transient DB
// outage does not silently under-recover. We probe the
// current row state via GetExecutionByID: a clean read
// with Status != "approved" confirms the race; any other
// outcome (read error, still-approved row) is a real
// failure worth propagating.
current, getErr := m.config.GetExecutionByID(ctx, exec.ExecutionID)
if getErr == nil && current != nil && current.Status != "approved" {
logging.Warnf("Skipping recovery of %s (already transitioned out of approved): %v", exec.ExecutionID, txErr)
continue

// AWS-only re-drive path (issue #632 Option 5): all AWS executors honour
// opts.IdempotencyToken via DeriveIdempotencyToken(exec.ExecutionID, i),
// so a second call with the same ExecutionID is a safe no-op on the AWS
// side. The ExecutionID must be non-empty to derive a unique token; an
// empty ID would map every legacy row to the same token set.
if allRecsAWS(exec) && exec.ExecutionID != "" {
counted, driveErr := m.claimAndRedrive(ctx, exec)
if driveErr != nil {
return recovered, driveErr
}
return recovered, fmt.Errorf("failed to transition stranded execution %s to failed: %w", exec.ExecutionID, txErr)
if counted {
recovered++
}
continue
Comment thread
coderabbitai[bot] marked this conversation as resolved.
}

updated.Error = "execution was approved but its purchase run was interrupted before completing and never finalized; failed by the recovery sweep so it is not silently stuck (issue #632). Verify on the cloud provider that no commitment was created, then Retry."
if saveErr := m.config.SavePurchaseExecution(ctx, updated); saveErr != nil {
// The atomic flip to "failed" already landed via TransitionExecutionStatus;
// only the explanatory error string failed to persist. Log loudly but
// still count the recovery — the row is no longer stranded in "approved".
logging.Errorf("AUDIT GAP: failed to stamp recovery error on %s: %v", exec.ExecutionID, saveErr)
// Safe-fail path for mixed/Azure/GCP/legacy executions.
counted, failErr := m.safeFail(ctx, exec)
if failErr != nil {
return recovered, failErr
}
if counted {
recovered++
}
recovered++
}

return recovered, nil
Expand Down
Loading