Skip to content
Closed
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
32 changes: 25 additions & 7 deletions internal/server/handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -172,17 +172,35 @@ func (app *Application) handleCleanupExpiredRecords(ctx context.Context) (map[st
return result, nil
}

// handleReapStuckPurchases sweeps purchase_executions stuck in
// approved/running longer than the configured threshold and flips them
// to "failed" via the existing TransitionExecutionStatus CAS. See
// internal/purchase/reaper.go + issue #678 for the full rationale and
// safety properties.
// handleReapStuckPurchases recovers and reaps stranded purchase executions.
//
// The threshold is read fresh from the PURCHASE_APPROVED_REAP_AFTER env
// var on every invocation (not cached at startup) so an ops tune via
// It runs two sweeps in order:
// 1. RecoverStrandedApprovals (issue #632): for executions stranded in
// "approved" past the internal staleApprovedThreshold, idempotently
// re-drives AWS-only rows to completion so an interrupted sync run
// actually finishes, and safe-fails mixed/Azure/GCP/legacy rows.
// This must run BEFORE the reaper so AWS strands get a chance to
// complete rather than being unconditionally failed.
// 2. ReapStuckExecutions (issue #678): flips anything still stuck in
// approved/running past the configured threshold to "failed" via the
// existing TransitionExecutionStatus CAS, so no row can be permanently
// stranded. See internal/purchase/reaper.go for the safety properties.
//
// The reap threshold is read fresh from the PURCHASE_APPROVED_REAP_AFTER
// env var on every invocation (not cached at startup) so an ops tune via
// Lambda env-var rotation takes effect on the next sweep without a
// redeploy.
//
// A recovery-sweep error is logged but NOT propagated: it must not block
// the reaper, which is the durable safety net that guarantees stranded
// rows always reach a terminal, retry-able state.
func (app *Application) handleReapStuckPurchases(ctx context.Context) (*purchase.ReapResult, error) {
if recovered, recErr := app.Purchase.RecoverStrandedApprovals(ctx); recErr != nil {
log.Printf("Stranded-approval recovery sweep error (continuing to reaper): %v", recErr)
} else if recovered > 0 {
log.Printf("Stranded-approval recovery sweep recovered %d execution(s)", recovered)
}

reapAfter := purchase.ParseReapAfterFromEnv()
log.Printf("Reaping stuck purchase executions (threshold: %s)...", reapAfter)
result, err := app.Purchase.ReapStuckExecutions(ctx, reapAfter)
Expand Down
59 changes: 59 additions & 0 deletions internal/server/handler_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,65 @@ func TestHandleScheduledTask(t *testing.T) {
}
}

// TestHandleReapStuckPurchasesRunsRecoveryFirst is the regression test for
// issue #632: the AWS-scoped re-drive sweep (RecoverStrandedApprovals, added
// by PR #728) was built and unit-tested but never dispatched at runtime, so
// stranded "approved" executions were only ever failed by the reaper instead
// of being re-driven to completion. This asserts that the reap_stuck_purchases
// task now invokes RecoverStrandedApprovals BEFORE ReapStuckExecutions, and
// that a recovery error does not block the reaper (the durable safety net).
func TestHandleReapStuckPurchasesRunsRecoveryFirst(t *testing.T) {
t.Run("recovery runs before reaper", func(t *testing.T) {
t.Setenv("PURCHASE_APPROVED_REAP_AFTER", "")
ctx := testutil.TestContext(t)

var order []string
mockPurchase := &testutil.MockPurchaseManager{
RecoverStrandedApprovalsFunc: func(ctx context.Context) (int, error) {
order = append(order, "recover")
return 1, nil
},
ReapStuckExecutionsFunc: func(ctx context.Context, reapAfter time.Duration) (*purchase.ReapResult, error) {
order = append(order, "reap")
return &purchase.ReapResult{}, nil
},
}

app := &Application{Scheduler: &testutil.MockScheduler{}, Purchase: mockPurchase}
_, err := app.HandleScheduledTask(ctx, TaskReapStuckPurchases)
testutil.AssertNoError(t, err)

if len(order) != 2 || order[0] != "recover" || order[1] != "reap" {
t.Fatalf("expected recovery then reaper, got %v", order)
}
})

t.Run("recovery error does not block reaper", func(t *testing.T) {
t.Setenv("PURCHASE_APPROVED_REAP_AFTER", "")
ctx := testutil.TestContext(t)

reaped := false
mockPurchase := &testutil.MockPurchaseManager{
RecoverStrandedApprovalsFunc: func(ctx context.Context) (int, error) {
return 0, errors.New("recovery sweep db error")
},
ReapStuckExecutionsFunc: func(ctx context.Context, reapAfter time.Duration) (*purchase.ReapResult, error) {
reaped = true
return &purchase.ReapResult{}, nil
},
}

app := &Application{Scheduler: &testutil.MockScheduler{}, Purchase: mockPurchase}
_, err := app.HandleScheduledTask(ctx, TaskReapStuckPurchases)
// Recovery failure is logged, not propagated; the reaper still runs
// and the task succeeds on the reaper's result.
testutil.AssertNoError(t, err)
if !reaped {
t.Fatal("reaper must still run when recovery sweep errors")
}
})
}

func TestTaskLockID(t *testing.T) {
t.Run("deterministic", func(t *testing.T) {
id1 := taskLockID(TaskCollectRecommendations)
Expand Down
9 changes: 9 additions & 0 deletions internal/server/interfaces.go
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,15 @@ type PurchaseManagerInterface interface {
ApproveExecution(ctx context.Context, execID, token, actor string) error
ApproveAndExecute(ctx context.Context, execID, actor string) error
CancelExecution(ctx context.Context, execID, token, actor string) error
// RecoverStrandedApprovals sweeps purchase_executions stranded in the
// "approved" status past the internal staleApprovedThreshold and, for
// AWS-only executions, idempotently re-drives them to completion (so an
// interrupted sync run actually finishes) before the reaper would fail
// them; mixed/Azure/GCP/legacy rows are safe-failed for visibility.
// Returns the number of rows acted on. Wired into the
// "reap_stuck_purchases" scheduled task ahead of ReapStuckExecutions.
// See issue #632.
RecoverStrandedApprovals(ctx context.Context) (int, error)
// ReapStuckExecutions sweeps purchase_executions stuck in
// approved/running longer than reapAfter and flips them to "failed"
// via the existing TransitionExecutionStatus CAS. Wired into the
Expand Down
8 changes: 8 additions & 0 deletions internal/testutil/mocks.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@ type MockPurchaseManager struct {
ApproveExecutionFunc func(ctx context.Context, execID, token, actor string) error
ApproveAndExecuteFunc func(ctx context.Context, execID, actor string) error
CancelExecutionFunc func(ctx context.Context, execID, token, actor string) error
RecoverStrandedApprovalsFunc func(ctx context.Context) (int, error)
ReapStuckExecutionsFunc func(ctx context.Context, reapAfter time.Duration) (*purchase.ReapResult, error)
}

Expand Down Expand Up @@ -90,6 +91,13 @@ func (m *MockPurchaseManager) CancelExecution(ctx context.Context, execID, token
return nil
}

func (m *MockPurchaseManager) RecoverStrandedApprovals(ctx context.Context) (int, error) {
if m.RecoverStrandedApprovalsFunc != nil {
return m.RecoverStrandedApprovalsFunc(ctx)
}
return 0, nil
}

func (m *MockPurchaseManager) ReapStuckExecutions(ctx context.Context, reapAfter time.Duration) (*purchase.ReapResult, error) {
if m.ReapStuckExecutionsFunc != nil {
return m.ReapStuckExecutionsFunc(ctx, reapAfter)
Expand Down
Loading