Skip to content
12 changes: 12 additions & 0 deletions internal/api/handler_purchases.go
Original file line number Diff line number Diff line change
Expand Up @@ -1064,12 +1064,24 @@ func (h *Handler) persistRetryExecution(ctx context.Context, failedExec *config.
tokenExpiresAt := time.Now().Add(config.ApprovalTokenTTL)

newExecutionID := uuid.New().String()
// Carry the predecessor's stable idempotency lineage key onto the retry
// successor VERBATIM (issue #1012) so the per-rec provider token is
// reproduced and a re-drive of a "failed" execution whose commitment
// actually landed short-circuits at the provider instead of double-buying.
// Legacy predecessors (persisted before migration 000066) carry no key —
// seed the successor's key from the predecessor's ExecutionID, which is what
// the original attempt's token derived from, so the match still holds.
idempotencyKey := failedExec.IdempotencyKey
if idempotencyKey == "" {
idempotencyKey = failedExec.ExecutionID
}
// PlanID + StepNumber propagate from the predecessor so a retried
// planned execution stays attributed to its plan + ramp step (CR
// #168 review). For ad-hoc executions PlanID is "" and StepNumber
// is 0, so propagation is a no-op for the non-plan case.
newExecution := &config.PurchaseExecution{
ExecutionID: newExecutionID,
IdempotencyKey: idempotencyKey,
PlanID: failedExec.PlanID,
StepNumber: failedExec.StepNumber,
Status: "pending",
Expand Down
12 changes: 9 additions & 3 deletions internal/commitmentopts/service.go
Original file line number Diff line number Diff line change
Expand Up @@ -162,9 +162,15 @@ func (s *Service) Validate(ctx context.Context, provider, service string, term i
}
combos, ok := byService[service]
if !ok {
// We have data for this provider but not this service. The
// service is not commitment-capable in our probe set (e.g.
// Savings Plans) — permissive fallback.
// We have probe data for this provider but not for this service.
// This typically means the service is not commitment-capable per
// the probe (e.g. Savings Plans has no per-service offering list)
// or it is a new service not yet covered by the probe set.
// Log at Warn so operators can detect misconfigured service names
// or probe gaps without blocking the plan save (05-M4). The
// frontend's hardcoded rules are the primary user-facing gate.
logging.Warnf("commitmentopts: provider %q known in probe data but service %q absent; treating as valid (term=%d payment=%q)",
provider, service, term, payment)
return true, nil
}
for _, c := range combos {
Expand Down
90 changes: 63 additions & 27 deletions internal/config/store_postgres.go
Original file line number Diff line number Diff line change
Expand Up @@ -723,6 +723,16 @@ func (s *PostgresStore) SavePurchaseExecutionTx(ctx context.Context, tx pgx.Tx,
execution.ExecutionID = uuid.New().String()
}

// Stamp a stable idempotency lineage key on first creation (issue #1012).
// Unlike ExecutionID this is generated once and then copied verbatim onto
// Retry successors / multi-account fan-out rows by the callers, so the
// derived provider token survives a re-drive. INSERT-only below (omitted
// from the ON CONFLICT DO UPDATE SET), so an upsert of an existing row
// never overwrites the key already persisted at creation.
if execution.IdempotencyKey == "" {
execution.IdempotencyKey = uuid.New().String()
}

// Marshal recommendations to JSONB
recommendationsJSON, err := json.Marshal(execution.Recommendations)
if err != nil {
Expand All @@ -744,8 +754,9 @@ func (s *PostgresStore) SavePurchaseExecutionTx(ctx context.Context, tx pgx.Tx,
cloud_account_id, source, approved_by, cancelled_by, capacity_percent,
created_by_user_id, retry_execution_id, retry_attempt_n,
approval_token_expires_at,
executed_by_user_id, executed_at, pre_approval_skip_reason
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $25)
executed_by_user_id, executed_at, pre_approval_skip_reason,
idempotency_key
) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17, $18, $19, $20, $21, $22, $23, $24, $25, $26)
ON CONFLICT (execution_id) DO UPDATE SET
status = $3,
notification_sent = $6,
Expand Down Expand Up @@ -815,6 +826,7 @@ func (s *PostgresStore) SavePurchaseExecutionTx(ctx context.Context, tx pgx.Tx,
execution.ExecutedByUserID,
execution.ExecutedAt,
execution.PreApprovalSkipReason,
execution.IdempotencyKey,
)

if err != nil {
Expand All @@ -838,7 +850,8 @@ func (s *PostgresStore) TransitionExecutionStatus(ctx context.Context, execution
cloud_account_id, source, approved_by, cancelled_by, capacity_percent,
created_by_user_id, retry_execution_id, retry_attempt_n,
approval_token_expires_at,
executed_by_user_id, executed_at, pre_approval_skip_reason
executed_by_user_id, executed_at, pre_approval_skip_reason,
idempotency_key
`

records, err := s.queryExecutions(ctx, query, executionID, toStatus, fromStatuses)
Expand Down Expand Up @@ -942,7 +955,8 @@ func (s *PostgresStore) GetExecutionsByStatuses(ctx context.Context, statuses []
cloud_account_id, source, approved_by, cancelled_by, capacity_percent,
created_by_user_id, retry_execution_id, retry_attempt_n,
approval_token_expires_at,
executed_by_user_id, executed_at, pre_approval_skip_reason
executed_by_user_id, executed_at, pre_approval_skip_reason,
idempotency_key
FROM purchase_executions
WHERE status = ANY($1)
ORDER BY scheduled_date DESC
Expand Down Expand Up @@ -981,7 +995,9 @@ func (s *PostgresStore) GetPlannedExecutions(ctx context.Context, statuses []str
total_upfront_cost, estimated_savings, completed_at, error, expires_at,
cloud_account_id, source, approved_by, cancelled_by, capacity_percent,
created_by_user_id, retry_execution_id, retry_attempt_n,
approval_token_expires_at
approval_token_expires_at,
executed_by_user_id, executed_at, pre_approval_skip_reason,
idempotency_key
FROM purchase_executions
WHERE status = ANY($1)
ORDER BY scheduled_date ASC NULLS LAST, id ASC
Expand All @@ -1005,7 +1021,8 @@ func (s *PostgresStore) GetStaleApprovedExecutions(ctx context.Context, olderTha
cloud_account_id, source, approved_by, cancelled_by, capacity_percent,
created_by_user_id, retry_execution_id, retry_attempt_n,
approval_token_expires_at,
executed_by_user_id, executed_at, pre_approval_skip_reason
executed_by_user_id, executed_at, pre_approval_skip_reason,
idempotency_key
FROM purchase_executions
WHERE status = 'approved' AND updated_at < NOW() - $1::interval
`
Expand Down Expand Up @@ -1046,7 +1063,8 @@ func (s *PostgresStore) ListStuckExecutions(ctx context.Context, statuses []stri
cloud_account_id, source, approved_by, cancelled_by, capacity_percent,
created_by_user_id, retry_execution_id, retry_attempt_n,
approval_token_expires_at,
executed_by_user_id, executed_at, pre_approval_skip_reason
executed_by_user_id, executed_at, pre_approval_skip_reason,
idempotency_key
FROM purchase_executions
WHERE status = ANY($1)
AND updated_at < NOW() - $2::interval
Expand All @@ -1066,7 +1084,8 @@ func (s *PostgresStore) GetPendingExecutions(ctx context.Context) ([]PurchaseExe
cloud_account_id, source, approved_by, cancelled_by, capacity_percent,
created_by_user_id, retry_execution_id, retry_attempt_n,
approval_token_expires_at,
executed_by_user_id, executed_at, pre_approval_skip_reason
executed_by_user_id, executed_at, pre_approval_skip_reason,
idempotency_key
FROM purchase_executions
WHERE status IN ('pending', 'notified')
AND (expires_at IS NULL OR expires_at > NOW())
Expand All @@ -1089,7 +1108,8 @@ func (s *PostgresStore) GetPendingExecutionsTx(ctx context.Context, tx pgx.Tx) (
cloud_account_id, source, approved_by, cancelled_by, capacity_percent,
created_by_user_id, retry_execution_id, retry_attempt_n,
approval_token_expires_at,
executed_by_user_id, executed_at, pre_approval_skip_reason
executed_by_user_id, executed_at, pre_approval_skip_reason,
idempotency_key
FROM purchase_executions
WHERE status IN ('pending', 'notified')
AND (expires_at IS NULL OR expires_at > NOW())
Expand All @@ -1114,7 +1134,8 @@ func (s *PostgresStore) GetExecutionByID(ctx context.Context, executionID string
cloud_account_id, source, approved_by, cancelled_by, capacity_percent,
created_by_user_id, retry_execution_id, retry_attempt_n,
approval_token_expires_at,
executed_by_user_id, executed_at, pre_approval_skip_reason
executed_by_user_id, executed_at, pre_approval_skip_reason,
idempotency_key
FROM purchase_executions
WHERE execution_id = $1
`
Expand All @@ -1140,7 +1161,8 @@ func (s *PostgresStore) GetExecutionByPlanAndDate(ctx context.Context, planID st
cloud_account_id, source, approved_by, cancelled_by, capacity_percent,
created_by_user_id, retry_execution_id, retry_attempt_n,
approval_token_expires_at,
executed_by_user_id, executed_at, pre_approval_skip_reason
executed_by_user_id, executed_at, pre_approval_skip_reason,
idempotency_key
FROM purchase_executions
WHERE plan_id = $1 AND scheduled_date = $2
`
Expand Down Expand Up @@ -1228,6 +1250,10 @@ func scanExecutionRows(rows pgx.Rows) ([]PurchaseExecution, error) {
// plan_id is nullable since migration 000033 (direct-execute
// rows from the Recommendations page have no originating plan).
var planID sql.NullString
// idempotency_key is NULL on rows created before migration 000066;
// leave exec.IdempotencyKey "" for those so the derivation falls back
// to ExecutionID (issue #1012).
var idempotencyKey sql.NullString

err := rows.Scan(
&planID,
Expand Down Expand Up @@ -1255,6 +1281,7 @@ func scanExecutionRows(rows pgx.Rows) ([]PurchaseExecution, error) {
&exec.ExecutedByUserID,
&executedAt,
&exec.PreApprovalSkipReason,
&idempotencyKey,
)
if err != nil {
return nil, fmt.Errorf("failed to scan execution: %w", err)
Expand All @@ -1263,35 +1290,44 @@ func scanExecutionRows(rows pgx.Rows) ([]PurchaseExecution, error) {
if planID.Valid {
exec.PlanID = planID.String
}
if idempotencyKey.Valid {
exec.IdempotencyKey = idempotencyKey.String
}

// Unmarshal recommendations
if err := json.Unmarshal(recommendationsJSON, &exec.Recommendations); err != nil {
return nil, fmt.Errorf("failed to unmarshal recommendations: %w", err)
}

// Handle nullable timestamps
if notifSent.Valid {
exec.NotificationSent = &notifSent.Time
}
if completedAt.Valid {
exec.CompletedAt = &completedAt.Time
}
if expiresAt.Valid {
exec.TTL = ttlFromTime(expiresAt.Time)
}
if tokenExpiresAt.Valid {
exec.ApprovalTokenExpiresAt = &tokenExpiresAt.Time
}
if executedAt.Valid {
exec.ExecutedAt = &executedAt.Time
}
applyExecutionNullableTimes(&exec, notifSent, completedAt, expiresAt, tokenExpiresAt, executedAt)

executions = append(executions, exec)
}

return executions, rows.Err()
}

// applyExecutionNullableTimes maps the nullable timestamp columns onto exec.
// It pulls the per-field NULL handling out of scanExecutionRows to keep that
// function under the cyclomatic limit.
func applyExecutionNullableTimes(exec *PurchaseExecution, notifSent, completedAt, expiresAt, tokenExpiresAt, executedAt sql.NullTime) {
if notifSent.Valid {
exec.NotificationSent = &notifSent.Time
}
if completedAt.Valid {
exec.CompletedAt = &completedAt.Time
}
if expiresAt.Valid {
exec.TTL = ttlFromTime(expiresAt.Time)
}
if tokenExpiresAt.Valid {
exec.ApprovalTokenExpiresAt = &tokenExpiresAt.Time
}
if executedAt.Valid {
exec.ExecutedAt = &executedAt.Time
}
}

// CleanupOldExecutions deletes purchase executions older than retentionDays.
//
// Two independent cleanup branches, each with its own retention window so
Expand Down
65 changes: 65 additions & 0 deletions internal/config/store_postgres_pgxmock_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -436,6 +436,7 @@ func TestPGXMock_GetExecutionByID_Success(t *testing.T) {
"created_by_user_id", "retry_execution_id", "retry_attempt_n",
"approval_token_expires_at",
"executed_by_user_id", "executed_at", "pre_approval_skip_reason",
"idempotency_key",
}
rows := pgxmock.NewRows(cols).AddRow(
"plan-1", "exec-1", "pending", 1, now,
Expand All @@ -445,6 +446,7 @@ func TestPGXMock_GetExecutionByID_Success(t *testing.T) {
nil, nil, 0,
sql.NullTime{},
nil, sql.NullTime{}, nil,
nil, // idempotency_key (NULL: legacy-row scan path, migration 000066)
)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
mock.ExpectQuery("SELECT").WithArgs(pgxmock.AnyArg()).WillReturnRows(rows)

Expand All @@ -459,6 +461,10 @@ func TestPGXMock_GetExecutionByID_Success(t *testing.T) {
// columns will catch.
assert.Nil(t, exec.RetryExecutionID)
assert.Equal(t, 0, exec.RetryAttemptN)
// idempotency_key scan-outcome guard (CR): the NULL path (legacy rows
// before migration 000066) must leave IdempotencyKey empty so derivation
// falls back to ExecutionID (issue #1012).
assert.Equal(t, "", exec.IdempotencyKey)
assert.NoError(t, mock.ExpectationsWereMet())
}

Expand All @@ -478,6 +484,7 @@ func TestPGXMock_GetExecutionByID_WithTimestamps(t *testing.T) {
"created_by_user_id", "retry_execution_id", "retry_attempt_n",
"approval_token_expires_at",
"executed_by_user_id", "executed_at", "pre_approval_skip_reason",
"idempotency_key",
}
successorID := "exec-3"
rows := pgxmock.NewRows(cols).AddRow(
Expand All @@ -491,6 +498,7 @@ func TestPGXMock_GetExecutionByID_WithTimestamps(t *testing.T) {
nil, &successorID, 2,
sql.NullTime{},
nil, sql.NullTime{}, nil,
"idem-key-exec-2", // idempotency_key non-NULL: exercises the scan path
)
mock.ExpectQuery("SELECT").WithArgs(pgxmock.AnyArg()).WillReturnRows(rows)

Expand All @@ -507,6 +515,61 @@ func TestPGXMock_GetExecutionByID_WithTimestamps(t *testing.T) {
require.NotNil(t, exec.RetryExecutionID)
assert.Equal(t, "exec-3", *exec.RetryExecutionID)
assert.Equal(t, 2, exec.RetryAttemptN)
// idempotency_key scan-outcome guard (CR): the non-NULL path must scan
// the stored key into IdempotencyKey verbatim.
assert.Equal(t, "idem-key-exec-2", exec.IdempotencyKey)
assert.NoError(t, mock.ExpectationsWereMet())
}

// TestPGXMock_GetPlannedExecutions_ProjectsAllScanColumns guards against the
// SELECT projection in GetPlannedExecutions drifting out of sync with the
// scanExecutionRows Scan target. When migration 000066 added idempotency_key
// (and earlier migrations added executed_by_user_id, executed_at,
// pre_approval_skip_reason) every execution-reading SELECT must project them,
// or scanExecutionRows fails at runtime with "failed to scan execution" on the
// planned-purchase list path (handler_purchases.go GetPlannedExecutions).
//
// The mock uses regexp query matching, so ExpectQuery requires the issued SQL
// to contain idempotency_key; with the column missing from the projection the
// query does not match and GetPlannedExecutions returns an error, which is what
// this test asserts against. It fails on the pre-fix projection and passes once
// the four trailing columns are added.
func TestPGXMock_GetPlannedExecutions_ProjectsAllScanColumns(t *testing.T) {
mock := newMock(t)
store := storeWith(mock)
ctx := context.Background()

recsJSON, _ := json.Marshal([]RecommendationRecord{})
now := time.Now().Truncate(time.Second)
cols := []string{
"plan_id", "execution_id", "status", "step_number", "scheduled_date",
"notification_sent", "approval_token", "recommendations",
"total_upfront_cost", "estimated_savings", "completed_at", "error", "expires_at",
"cloud_account_id", "source", "approved_by", "cancelled_by", "capacity_percent",
"created_by_user_id", "retry_execution_id", "retry_attempt_n",
"approval_token_expires_at",
"executed_by_user_id", "executed_at", "pre_approval_skip_reason",
"idempotency_key",
}
rows := pgxmock.NewRows(cols).AddRow(
"plan-1", "exec-1", "pending", 1, now,
sql.NullTime{}, "tok-123", recsJSON,
100.0, 200.0, sql.NullTime{}, "", sql.NullTime{},
nil, "", nil, nil, 100,
nil, nil, 0,
sql.NullTime{},
nil, sql.NullTime{}, nil,
"idem-key-planned",
)
// Regexp matcher: only matches if the issued SELECT projects idempotency_key.
mock.ExpectQuery("idempotency_key").
WithArgs(pgxmock.AnyArg(), pgxmock.AnyArg()).
WillReturnRows(rows)

execs, err := store.GetPlannedExecutions(ctx, []string{"pending"}, 10)
require.NoError(t, err)
require.Len(t, execs, 1)
assert.Equal(t, "idem-key-planned", execs[0].IdempotencyKey)
assert.NoError(t, mock.ExpectationsWereMet())
}

Expand Down Expand Up @@ -1800,6 +1863,7 @@ func stuckExecRow(execID, status string, scheduled time.Time) []any {
nil, nil, 0,
sql.NullTime{},
nil, sql.NullTime{}, nil,
nil, // idempotency_key (NULL: legacy-row scan path, migration 000066)
}
}

Expand All @@ -1812,6 +1876,7 @@ func stuckExecCols() []string {
"created_by_user_id", "retry_execution_id", "retry_attempt_n",
"approval_token_expires_at",
"executed_by_user_id", "executed_at", "pre_approval_skip_reason",
"idempotency_key",
}
}

Expand Down
11 changes: 11 additions & 0 deletions internal/config/types.go
Original file line number Diff line number Diff line change
Expand Up @@ -286,6 +286,17 @@ type PurchaseExecution struct {
// string "direct-execute permission". NULL on every normal-flow row.
// Migration 000058.
PreApprovalSkipReason *string `json:"pre_approval_skip_reason,omitempty" dynamodbav:"pre_approval_skip_reason,omitempty"`
// IdempotencyKey is the stable lineage anchor the per-rec provider
// idempotency token is derived from (issue #1012). Unlike ExecutionID
// it is NOT regenerated on Retry or multi-account fan-out: it is
// generated once at first creation, copied verbatim onto every Retry
// successor, and combined with the account ID to seed each per-account
// fan-out row. This makes DeriveIdempotencyToken reproduce the same
// token across a strand-and-re-drive so the provider dedupes and the
// commitment is never bought twice. Empty on rows created before
// migration 000066 — the derivation falls back to ExecutionID for those
// (identical to the pre-fix behaviour for a single un-retried execution).
IdempotencyKey string `json:"idempotency_key,omitempty" dynamodbav:"idempotency_key,omitempty"`
}

// IsCancelable reports whether an execution may still be cancelled. Only the
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
-- Reverse 000066: drop the idempotency lineage key column.
ALTER TABLE purchase_executions
DROP COLUMN IF EXISTS idempotency_key;
Loading
Loading